diff --git a/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbLegacyMigrator.cs b/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbLegacyMigrator.cs
index 42bf6e42..29f0eddd 100644
--- a/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbLegacyMigrator.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbLegacyMigrator.cs
@@ -83,10 +83,16 @@ public static class SiteLocalDbLegacyMigrator
///
/// notification_lists and smtp_configurations are deliberately absent.
/// Both are purged on every deploy and are permanently empty by design — the site-side
- /// write paths were removed on 2026-07-10. A pre-fix legacy file can still hold rows,
- /// and smtp_configurations.password is plaintext, so migrating them would
- /// resurrect plaintext SMTP passwords into a table that Task 14 makes REPLICATED. The
- /// tables themselves are still created (see SiteStorageSchema); only their
+ /// write paths were removed on 2026-07-10. A pre-fix legacy file can still hold rows, and
+ /// smtp_configurations.password is plaintext.
+ ///
+ /// Skipping them here is one half of a pair: the cutover also declines to register them
+ /// for replication, for the same reason. Migrating them would leave plaintext SMTP
+ /// passwords sitting in the consolidated database — one future RegisterReplicated
+ /// away from being shipped to a peer — in exchange for resurrecting config that nothing
+ /// reads. Keeping the tables permanently empty is what makes both decisions safe.
+ ///
+ /// The tables themselves are still created (see SiteStorageSchema); only their
/// historical contents are left behind.
///
internal static readonly LegacyTable[] SiteStorageTables =
diff --git a/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbLegacyMigratorTests.cs b/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbLegacyMigratorTests.cs
index 834aa36a..329b1011 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbLegacyMigratorTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbLegacyMigratorTests.cs
@@ -465,9 +465,10 @@ public class SiteLocalDbLegacyMigratorTests : IDisposable
{
// smtp_configurations.password is PLAINTEXT. Both tables are purged on every deploy
// and permanently empty by design since the site write paths were removed
- // (2026-07-10), but a pre-fix legacy file can still hold rows — and Task 14 makes
- // these tables REPLICATED. Migrating them would push plaintext SMTP passwords across
- // a replication channel whose only historical payload was exactly that.
+ // (2026-07-10), but a pre-fix legacy file can still hold rows. The cutover also
+ // declines to REGISTER these two for replication, for the same reason — so migrating
+ // them would leave plaintext SMTP passwords in the consolidated database, one future
+ // RegisterReplicated away from being shipped to a peer, for config nothing reads.
var sitePath = Path_("scadabridge.db");
SeedLegacySiteStorage(sitePath);
var config = Config(siteDbPath: sitePath);
diff --git a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairConvergenceTests.cs b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairConvergenceTests.cs
index b9c13ced..a8ee9ad1 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairConvergenceTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairConvergenceTests.cs
@@ -1,202 +1,23 @@
-using System.Net;
-using Microsoft.AspNetCore.Builder;
-using Microsoft.AspNetCore.Hosting;
-using Microsoft.AspNetCore.Hosting.Server;
-using Microsoft.AspNetCore.Hosting.Server.Features;
-using Microsoft.AspNetCore.Server.Kestrel.Core;
-using Microsoft.Extensions.Configuration;
-using Microsoft.Extensions.DependencyInjection;
-using Microsoft.Extensions.Hosting;
using ZB.MOM.WW.LocalDb;
-using ZB.MOM.WW.LocalDb.Replication;
-using ZB.MOM.WW.ScadaBridge.Host;
namespace ZB.MOM.WW.ScadaBridge.IntegrationTests;
///
-/// Serializes the site-pair convergence tests against each other: each one stands up a real
-/// Kestrel listener plus two SQLite files, and running them concurrently under CI
-/// contention is a flakiness risk.
-///
-[CollectionDefinition("LocalDbSitePairConvergence")]
-public sealed class LocalDbSitePairConvergenceCollection;
-
-///
-/// Two ScadaBridge site nodes replicating the consolidated site database over a REAL
-/// loopback gRPC transport, through the REAL fail-closed auth interceptor.
+/// Phase 1 convergence: operation tracking and site events across a real site pair.
///
///
-///
/// This is the test that answers the question Phase 1 exists to answer: does a site node
/// pair actually stop losing operation-tracking and site-event state? Everything upstream
/// of it — schema helpers, DI wiring, the interceptor — can be individually green while the
/// pair still fails to converge.
-///
///
-/// It uses , not a hand-written schema, so the tables,
-/// their primary keys, and the registration ORDER under test are the ones the host actually
-/// runs. A separate schema here would prove only that the test agrees with itself.
-///
-///
-/// Offline: no docker, no external services. Loopback Kestrel with h2c.
+/// The fixture lives in , shared with the Phase 2
+/// convergence suites.
///
///
[Collection("LocalDbSitePairConvergence")]
-public sealed class LocalDbSitePairConvergenceTests : IAsyncLifetime
+public sealed class LocalDbSitePairConvergenceTests : LocalDbSitePairHarness
{
- private const string SharedApiKey = "site-pair-convergence-key";
- private static readonly TimeSpan ConvergeTimeout = TimeSpan.FromSeconds(30);
-
- private readonly string _pathA = Path.Combine(Path.GetTempPath(), $"sitepairA-{Guid.NewGuid():N}.db");
- private readonly string _pathB = Path.Combine(Path.GetTempPath(), $"sitepairB-{Guid.NewGuid():N}.db");
-
- // The databases are owned by the fixture, in their own providers, and registered into the
- // hosts as pre-constructed instances. MS.DI does not dispose instances it did not create,
- // so tearing a host down (the offline-peer scenario) leaves the databases intact and
- // writable — which is exactly what lets node A accumulate writes while B is down.
- private ServiceProvider _dbProviderA = null!;
- private ServiceProvider _dbProviderB = null!;
-
- private IHost? _serverHost; // node B — passive
- private IHost? _initiatorHost; // node A — dials the peer
-
- static LocalDbSitePairConvergenceTests() =>
- // Grpc.Net.Client dials the loopback server over HTTP/2 cleartext.
- AppContext.SetSwitch("System.Net.Http.SocketsHttpHandler.Http2UnencryptedSupport", true);
-
- private ILocalDb A => _dbProviderA.GetRequiredService();
- private ILocalDb B => _dbProviderB.GetRequiredService();
-
- public async Task InitializeAsync()
- {
- _dbProviderA = BuildDatabaseProvider(_pathA, "node-a");
- _dbProviderB = BuildDatabaseProvider(_pathB, "node-b");
-
- // Force construction (and therefore OnReady) before anything replicates.
- _ = A;
- _ = B;
-
- await StartPassiveAsync();
- await StartInitiatorAsync();
- }
-
- public async Task DisposeAsync()
- {
- await StopHostAsync(_initiatorHost);
- await StopHostAsync(_serverHost);
- await _dbProviderA.DisposeAsync();
- await _dbProviderB.DisposeAsync();
-
- Microsoft.Data.Sqlite.SqliteConnection.ClearAllPools();
- foreach (var path in new[] { _pathA, _pathB })
- {
- foreach (var suffix in new[] { "", "-wal", "-shm" })
- {
- try { File.Delete(path + suffix); } catch { /* best effort */ }
- }
- }
- }
-
- // ---- fixture internals ------------------------------------------------------------
-
- ///
- /// A provider owning one consolidated site database, initialized through the host's own
- /// — same schema, same registration order.
- ///
- private static ServiceProvider BuildDatabaseProvider(string path, string nodeName)
- {
- var config = new ConfigurationBuilder()
- .AddInMemoryCollection(new Dictionary
- {
- ["LocalDb:Path"] = path,
- ["ScadaBridge:Node:NodeName"] = nodeName,
- // Point the legacy migrator at paths that do not exist, so it no-ops rather
- // than picking up stray files from the test working directory.
- ["ScadaBridge:OperationTracking:ConnectionString"] =
- $"Data Source={Path.Combine(Path.GetTempPath(), $"absent-{Guid.NewGuid():N}.db")}",
- ["ScadaBridge:SiteEventLog:DatabasePath"] =
- Path.Combine(Path.GetTempPath(), $"absent-{Guid.NewGuid():N}.db"),
- })
- .Build();
-
- return new ServiceCollection()
- .AddZbLocalDb(config, db => SiteLocalDbSetup.OnReady(db, config))
- .BuildServiceProvider();
- }
-
- private static IConfiguration ReplicationConfig(string? peerAddress)
- {
- var values = new Dictionary
- {
- // Tight flush + bounded reconnect backoff so convergence is observable well
- // inside the poll deadline. The 60 s production default would let the doubling
- // backoff overrun it after a peer outage.
- ["LocalDb:Replication:FlushInterval"] = "00:00:00.050",
- ["LocalDb:Replication:ReconnectBackoffMax"] = "00:00:02",
- // Both nodes share one key — the interceptor is fail-closed, so a mismatch here
- // turns every scenario below red (verified by deliberately breaking it).
- ["LocalDb:Replication:ApiKey"] = SharedApiKey,
- };
-
- if (peerAddress is not null)
- values["LocalDb:Replication:PeerAddress"] = peerAddress;
-
- return new ConfigurationBuilder().AddInMemoryCollection(values).Build();
- }
-
- private async Task StartPassiveAsync()
- {
- var config = ReplicationConfig(peerAddress: null);
-
- _serverHost = await new HostBuilder()
- .ConfigureWebHost(web =>
- {
- web.UseKestrel(o =>
- o.Listen(IPAddress.Loopback, 0, listen => listen.Protocols = HttpProtocols.Http2));
- web.ConfigureServices(services =>
- {
- services.AddLogging();
- services.AddRouting();
- // The REAL interceptor, not a stand-in. If it rejected legitimate peer
- // traffic, every scenario below would fail — which is the point.
- services.AddGrpc(o => o.Interceptors.Add());
- services.AddSingleton(B);
- services.AddZbLocalDbReplication(config);
- });
- web.Configure(app =>
- {
- app.UseRouting();
- app.UseEndpoints(e => e.MapZbLocalDbSync());
- });
- })
- .StartAsync();
- }
-
- private async Task StartInitiatorAsync()
- {
- var config = ReplicationConfig(PassiveAddress());
-
- _initiatorHost = await new HostBuilder()
- .ConfigureServices(services =>
- {
- services.AddLogging();
- services.AddSingleton(A);
- services.AddZbLocalDbReplication(config);
- })
- .StartAsync();
- }
-
- private string PassiveAddress()
- => _serverHost!.Services.GetRequiredService()
- .Features.Get()!.Addresses.Single();
-
- private static async Task StopHostAsync(IHost? host)
- {
- if (host is null) return;
- try { await host.StopAsync(TimeSpan.FromSeconds(5)); } catch { /* teardown */ }
- host.Dispose();
- }
-
// ---- data helpers -----------------------------------------------------------------
private static Task WriteTrackingAsync(ILocalDb db, string id, string status, string target)
@@ -232,19 +53,6 @@ public sealed class LocalDbSitePairConvergenceTests : IAsyncLifetime
private static Task> ReadEventIdsAsync(ILocalDb db)
=> db.QueryAsync("SELECT id FROM site_events ORDER BY id", static r => r.GetString(0));
- /// Polls until true or the deadline passes.
- private static async Task WaitUntilAsync(Func> condition, string because)
- {
- var deadline = DateTime.UtcNow + ConvergeTimeout;
- while (DateTime.UtcNow < deadline)
- {
- if (await condition()) return;
- await Task.Delay(50);
- }
-
- Assert.Fail($"Timed out after {ConvergeTimeout.TotalSeconds:0}s waiting for: {because}");
- }
-
// ---- scenarios --------------------------------------------------------------------
[Fact]
@@ -322,19 +130,15 @@ public sealed class LocalDbSitePairConvergenceTests : IAsyncLifetime
{
// The failover case that motivates Phase 1: one node is down while the other keeps
// working, and nothing written during the outage may be lost.
- await StopHostAsync(_serverHost);
- _serverHost = null;
+ await StopPassiveAsync();
var idsDuringOutage = Enumerable.Range(0, 5).Select(_ => Guid.NewGuid().ToString("N")).ToList();
foreach (var id in idsDuringOutage) await WriteEventAsync(A, id, "written while B was down");
await WriteTrackingAsync(A, "op-during-outage", "Delivered", "ERP.GetOrder");
// Node B's database survived the host teardown (pre-constructed instance), so this is
- // a genuine rejoin rather than a fresh node. It comes back on a NEW loopback port;
- // the initiator's channel factory re-reads the peer address on each reconnect.
- await StartPassiveAsync();
- await StopHostAsync(_initiatorHost);
- await StartInitiatorAsync();
+ // a genuine rejoin rather than a fresh node.
+ await RestartPairAsync();
await WaitUntilAsync(
async () =>
diff --git a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairHarness.cs b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairHarness.cs
new file mode 100644
index 00000000..dc37eefb
--- /dev/null
+++ b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairHarness.cs
@@ -0,0 +1,286 @@
+using System.Net;
+using Microsoft.AspNetCore.Builder;
+using Microsoft.AspNetCore.Hosting;
+using Microsoft.AspNetCore.Hosting.Server;
+using Microsoft.AspNetCore.Hosting.Server.Features;
+using Microsoft.AspNetCore.Server.Kestrel.Core;
+using Microsoft.Extensions.Configuration;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Hosting;
+using ZB.MOM.WW.LocalDb;
+using ZB.MOM.WW.LocalDb.Replication;
+using ZB.MOM.WW.ScadaBridge.Host;
+
+namespace ZB.MOM.WW.ScadaBridge.IntegrationTests;
+
+///
+/// Serializes the site-pair convergence tests against each other: each one stands up a real
+/// Kestrel listener plus two SQLite files, and running them concurrently under CI
+/// contention is a flakiness risk.
+///
+[CollectionDefinition("LocalDbSitePairConvergence")]
+public sealed class LocalDbSitePairConvergenceCollection;
+
+///
+/// Two ScadaBridge site nodes replicating the consolidated site database over a REAL
+/// loopback gRPC transport, through the REAL fail-closed auth interceptor.
+///
+///
+///
+/// Extracted from the Phase 1 convergence tests once Phase 2 needed the same pair for the
+/// store-and-forward buffer and the configuration tables. Everything here is fixture; the
+/// derived classes hold only their own data helpers and scenarios.
+///
+///
+/// It uses , not a hand-written schema, so the tables,
+/// their primary keys, and the registration ORDER under test are the ones the host actually
+/// runs. A separate schema here would prove only that the test agrees with itself.
+///
+///
+/// Offline: no docker, no external services. Loopback Kestrel with h2c.
+///
+///
+public abstract class LocalDbSitePairHarness : IAsyncLifetime
+{
+ private const string SharedApiKey = "site-pair-convergence-key";
+
+ /// How long a scenario waits for the pair to agree before failing.
+ protected static readonly TimeSpan ConvergeTimeout = TimeSpan.FromSeconds(30);
+
+ private readonly string _pathA = Path.Combine(Path.GetTempPath(), $"sitepairA-{Guid.NewGuid():N}.db");
+ private readonly string _pathB = Path.Combine(Path.GetTempPath(), $"sitepairB-{Guid.NewGuid():N}.db");
+
+ // The databases are owned by the fixture, in their own providers, and registered into the
+ // hosts as pre-constructed instances. MS.DI does not dispose instances it did not create,
+ // so tearing a host down (the offline-peer scenario) leaves the databases intact and
+ // writable — which is exactly what lets node A accumulate writes while B is down.
+ private ServiceProvider _dbProviderA = null!;
+ private ServiceProvider _dbProviderB = null!;
+
+ private IHost? _serverHost; // node B — passive
+ private IHost? _initiatorHost; // node A — dials the peer
+
+ static LocalDbSitePairHarness() =>
+ // Grpc.Net.Client dials the loopback server over HTTP/2 cleartext.
+ AppContext.SetSwitch("System.Net.Http.SocketsHttpHandler.Http2UnencryptedSupport", true);
+
+ /// Node A — the initiator, which dials the peer.
+ protected ILocalDb A => _dbProviderA.GetRequiredService();
+
+ /// Node B — the passive node, which listens.
+ protected ILocalDb B => _dbProviderB.GetRequiredService();
+
+ public async Task InitializeAsync()
+ {
+ _dbProviderA = BuildDatabaseProvider(_pathA, "node-a");
+ _dbProviderB = BuildDatabaseProvider(_pathB, "node-b");
+
+ // Force construction (and therefore OnReady) before anything replicates.
+ _ = A;
+ _ = B;
+
+ await StartPassiveAsync();
+ await StartInitiatorAsync();
+ }
+
+ public async Task DisposeAsync()
+ {
+ await StopHostAsync(_initiatorHost);
+ await StopHostAsync(_serverHost);
+ await _dbProviderA.DisposeAsync();
+ await _dbProviderB.DisposeAsync();
+
+ Microsoft.Data.Sqlite.SqliteConnection.ClearAllPools();
+ foreach (var path in new[] { _pathA, _pathB })
+ {
+ foreach (var suffix in new[] { "", "-wal", "-shm" })
+ {
+ try { File.Delete(path + suffix); } catch { /* best effort */ }
+ }
+ }
+ }
+
+ // ---- fixture internals ------------------------------------------------------------
+
+ ///
+ /// A provider owning one consolidated site database, initialized through the host's own
+ /// — same schema, same registration order.
+ ///
+ private static ServiceProvider BuildDatabaseProvider(string path, string nodeName)
+ {
+ var config = new ConfigurationBuilder()
+ .AddInMemoryCollection(new Dictionary
+ {
+ ["LocalDb:Path"] = path,
+ ["ScadaBridge:Node:NodeName"] = nodeName,
+ // Point the legacy migrators at paths that do not exist, so they no-op rather
+ // than picking up stray files from the test working directory. The two Phase 2
+ // defaults matter most: unlike the Phase 1 pair they resolve inside ./data/,
+ // and a migration would also RENAME whatever it found.
+ ["ScadaBridge:OperationTracking:ConnectionString"] =
+ $"Data Source={Path.Combine(Path.GetTempPath(), $"absent-{Guid.NewGuid():N}.db")}",
+ ["ScadaBridge:SiteEventLog:DatabasePath"] =
+ Path.Combine(Path.GetTempPath(), $"absent-{Guid.NewGuid():N}.db"),
+ ["ScadaBridge:StoreAndForward:SqliteDbPath"] =
+ Path.Combine(Path.GetTempPath(), $"absent-{Guid.NewGuid():N}.db"),
+ ["ScadaBridge:Database:SiteDbPath"] =
+ Path.Combine(Path.GetTempPath(), $"absent-{Guid.NewGuid():N}.db"),
+ })
+ .Build();
+
+ return new ServiceCollection()
+ .AddZbLocalDb(config, db =>
+ {
+ SiteLocalDbSetup.OnReady(db, config);
+ RegisterPhase2TablesUntilCutover(db);
+ })
+ .BuildServiceProvider();
+ }
+
+ ///
+ /// Registers the Phase 2 tables for capture, which production OnReady does not do
+ /// until the Task 14 cutover.
+ ///
+ ///
+ /// Delete this method at Task 14, along with its call above. Until the cutover the
+ /// bespoke SiteReplicationActor and ReplicationService still own these
+ /// tables, so OnReady deliberately leaves them unregistered — but the convergence
+ /// specifications the cutover has to satisfy need capture triggers to mean anything. The
+ /// tables are empty at this point, so registering here captures nothing retroactively;
+ /// the ordering guarantee under test is unaffected.
+ ///
+ /// After Task 14 these registrations come from OnReady itself, and leaving this
+ /// method behind would mask a cutover that forgot one — the whole point of these tests.
+ ///
+ ///
+ private static void RegisterPhase2TablesUntilCutover(ILocalDb db)
+ {
+ foreach (var table in Phase2ReplicatedTables)
+ db.RegisterReplicated(table);
+ }
+
+ ///
+ /// The eight tables Task 14 registers. Listed literally rather than derived from
+ /// production code, so that a cutover which registers the wrong set fails these tests
+ /// instead of agreeing with itself.
+ ///
+ ///
+ /// notification_lists and smtp_configurations are absent by design. They
+ /// are permanently empty (no site writer since 2026-07-10, the migrator skips them, the
+ /// active-node purge keeps them empty), so registering them would open a standing
+ /// replication channel whose only historical payload was plaintext SMTP passwords.
+ ///
+ protected static readonly string[] Phase2ReplicatedTables =
+ [
+ "sf_messages",
+ "deployed_configurations", "static_attribute_overrides", "shared_scripts",
+ "external_systems", "database_connections", "data_connection_definitions",
+ "native_alarm_state",
+ ];
+
+ private static IConfiguration ReplicationConfig(string? peerAddress)
+ {
+ var values = new Dictionary
+ {
+ // Tight flush + bounded reconnect backoff so convergence is observable well
+ // inside the poll deadline. The 60 s production default would let the doubling
+ // backoff overrun it after a peer outage.
+ ["LocalDb:Replication:FlushInterval"] = "00:00:00.050",
+ ["LocalDb:Replication:ReconnectBackoffMax"] = "00:00:02",
+ // Both nodes share one key — the interceptor is fail-closed, so a mismatch here
+ // turns every scenario below red (verified by deliberately breaking it).
+ ["LocalDb:Replication:ApiKey"] = SharedApiKey,
+ };
+
+ if (peerAddress is not null)
+ values["LocalDb:Replication:PeerAddress"] = peerAddress;
+
+ return new ConfigurationBuilder().AddInMemoryCollection(values).Build();
+ }
+
+ /// Starts node B, the passive listener.
+ protected async Task StartPassiveAsync()
+ {
+ var config = ReplicationConfig(peerAddress: null);
+
+ _serverHost = await new HostBuilder()
+ .ConfigureWebHost(web =>
+ {
+ web.UseKestrel(o =>
+ o.Listen(IPAddress.Loopback, 0, listen => listen.Protocols = HttpProtocols.Http2));
+ web.ConfigureServices(services =>
+ {
+ services.AddLogging();
+ services.AddRouting();
+ // The REAL interceptor, not a stand-in. If it rejected legitimate peer
+ // traffic, every scenario below would fail — which is the point.
+ services.AddGrpc(o => o.Interceptors.Add());
+ services.AddSingleton(B);
+ services.AddZbLocalDbReplication(config);
+ });
+ web.Configure(app =>
+ {
+ app.UseRouting();
+ app.UseEndpoints(e => e.MapZbLocalDbSync());
+ });
+ })
+ .StartAsync();
+ }
+
+ /// Starts node A, which dials the passive node.
+ protected async Task StartInitiatorAsync()
+ {
+ var config = ReplicationConfig(PassiveAddress());
+
+ _initiatorHost = await new HostBuilder()
+ .ConfigureServices(services =>
+ {
+ services.AddLogging();
+ services.AddSingleton(A);
+ services.AddZbLocalDbReplication(config);
+ })
+ .StartAsync();
+ }
+
+ /// Takes node B's listener down, leaving its database intact and writable.
+ protected async Task StopPassiveAsync()
+ {
+ await StopHostAsync(_serverHost);
+ _serverHost = null;
+ }
+
+ ///
+ /// Brings node B back on a NEW loopback port and re-dials from A. The initiator's channel
+ /// factory re-reads the peer address on each reconnect, so this is a genuine rejoin.
+ ///
+ protected async Task RestartPairAsync()
+ {
+ await StartPassiveAsync();
+ await StopHostAsync(_initiatorHost);
+ await StartInitiatorAsync();
+ }
+
+ private string PassiveAddress()
+ => _serverHost!.Services.GetRequiredService()
+ .Features.Get()!.Addresses.Single();
+
+ private static async Task StopHostAsync(IHost? host)
+ {
+ if (host is null) return;
+ try { await host.StopAsync(TimeSpan.FromSeconds(5)); } catch { /* teardown */ }
+ host.Dispose();
+ }
+
+ /// Polls until true or the deadline passes.
+ protected static async Task WaitUntilAsync(Func> condition, string because)
+ {
+ var deadline = DateTime.UtcNow + ConvergeTimeout;
+ while (DateTime.UtcNow < deadline)
+ {
+ if (await condition()) return;
+ await Task.Delay(50);
+ }
+
+ Assert.Fail($"Timed out after {ConvergeTimeout.TotalSeconds:0}s waiting for: {because}");
+ }
+}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbStoreAndForwardConvergenceTests.cs b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbStoreAndForwardConvergenceTests.cs
new file mode 100644
index 00000000..409f23d8
--- /dev/null
+++ b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbStoreAndForwardConvergenceTests.cs
@@ -0,0 +1,221 @@
+using ZB.MOM.WW.LocalDb;
+using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
+
+namespace ZB.MOM.WW.ScadaBridge.IntegrationTests;
+
+///
+/// Phase 2 convergence: the store-and-forward buffer across a real site pair.
+///
+///
+///
+/// These are specifications, written before the deletion they justify. The bespoke
+/// ReplicationService (an explicit Add/Remove/Park/Requeue operation stream over Akka)
+/// is deleted at Task 14 and replaced by LocalDb's trigger-based change capture. Each test
+/// here is one behaviour the old mechanism provided, restated as an outcome the replacement
+/// must still deliver — so the cutover has something to be judged against other than "it
+/// compiles".
+///
+///
+/// They are written in terms of rows, not operations, which is the whole point:
+/// under CDC there is no Add or Park message to observe, only a row that must end up in the
+/// right state on both nodes.
+///
+///
+/// One intent is deliberately not ported:
+/// ReplicationServiceTests.ReplicationOperations_AreDispatchedInIssueOrder — 200
+/// interleaved operations dispatched synchronously, observed in strict issue order. That
+/// asserts the mechanism (inline fire-and-forget dispatch), not an outcome, and CDC
+/// capture is asynchronous and batched by construction, so no honest port exists. Its
+/// portable content is the ordering outcome: an add followed by a remove must never
+/// converge to present. That is
+/// . This paragraph exists so a
+/// future reader does not conclude the test was dropped by accident.
+///
+///
+[Collection("LocalDbSitePairConvergence")]
+public sealed class LocalDbStoreAndForwardConvergenceTests : LocalDbSitePairHarness
+{
+ // ---- data helpers -----------------------------------------------------------------
+
+ /// Buffers a message, as StoreAndForwardStorage.EnqueueAsync does.
+ private static Task EnqueueAsync(ILocalDb db, string id, string target = "ERP.GetOrder")
+ => db.ExecuteAsync(
+ """
+ INSERT INTO sf_messages (
+ id, category, target, payload_json, retry_count, max_retries,
+ retry_interval_ms, created_at, status)
+ VALUES (@id, 0, @target, '{"order":1}', 0, 50, 30000, @now, @pending)
+ ON CONFLICT(id) DO UPDATE SET
+ status = excluded.status,
+ retry_count = excluded.retry_count;
+ """,
+ new
+ {
+ id,
+ target,
+ now = DateTime.UtcNow.ToString("o"),
+ pending = (int)StoreAndForwardMessageStatus.Pending,
+ });
+
+ /// Deletes a delivered message, as the successful-retry path does.
+ private static Task DeleteAsync(ILocalDb db, string id)
+ => db.ExecuteAsync("DELETE FROM sf_messages WHERE id = @id;", new { id });
+
+ private static Task SetStatusAsync(ILocalDb db, string id, StoreAndForwardMessageStatus status, int retryCount)
+ => db.ExecuteAsync(
+ "UPDATE sf_messages SET status = @status, retry_count = @retryCount WHERE id = @id;",
+ new { id, status = (int)status, retryCount });
+
+ private static async Task<(int Status, int RetryCount)?> ReadAsync(ILocalDb db, string id)
+ {
+ var rows = await db.QueryAsync(
+ "SELECT status, retry_count FROM sf_messages WHERE id = @id",
+ static r => (r.GetInt32(0), r.GetInt32(1)), new { id });
+ return rows.Count == 0 ? null : rows[0];
+ }
+
+ private static async Task ExistsAsync(ILocalDb db, string id)
+ => await ReadAsync(db, id) is not null;
+
+ private static async Task HasStatusAsync(
+ ILocalDb db, string id, StoreAndForwardMessageStatus status)
+ => await ReadAsync(db, id) is { } row && row.Status == (int)status;
+
+ // ---- scenarios --------------------------------------------------------------------
+
+ [Fact]
+ public async Task BufferedMessage_MaterialisesOnThePeer()
+ {
+ // Was: BufferingAMessage_ReplicatesAnAddOperation.
+ await EnqueueAsync(A, "msg-add");
+
+ await WaitUntilAsync(
+ () => ExistsAsync(B, "msg-add"),
+ "a message buffered on node A to appear on node B");
+ }
+
+ [Fact]
+ public async Task DeliveredMessage_DisappearsFromThePeer()
+ {
+ // Was: SuccessfulRetry_ReplicatesARemoveOperation. Row deletes propagate as
+ // tombstones under CDC — the peer must not keep re-delivering a message that node A
+ // already succeeded on.
+ await EnqueueAsync(A, "msg-remove");
+ await WaitUntilAsync(() => ExistsAsync(B, "msg-remove"), "the message to arrive on B first");
+
+ await DeleteAsync(A, "msg-remove");
+
+ await WaitUntilAsync(
+ async () => !await ExistsAsync(B, "msg-remove"),
+ "the delivered message to be gone from node B");
+ }
+
+ [Fact]
+ public async Task ParkedMessage_ShowsAsParkedOnThePeer()
+ {
+ // Was: ParkedMessage_ReplicatesAParkOperation.
+ await EnqueueAsync(A, "msg-park");
+ await WaitUntilAsync(() => ExistsAsync(B, "msg-park"), "the message to arrive on B first");
+
+ await SetStatusAsync(A, "msg-park", StoreAndForwardMessageStatus.Parked, retryCount: 50);
+
+ await WaitUntilAsync(
+ () => HasStatusAsync(B, "msg-park", StoreAndForwardMessageStatus.Parked),
+ "the parked status to reach node B");
+ }
+
+ [Fact]
+ public async Task RequeuedMessage_ResetsStatusAndRetryCountOnThePeer()
+ {
+ // Was: RetryingAParkedMessage_ReplicatesARequeueOperation +
+ // ApplyReplicatedOperation_Requeue_MovesStandbyRowBackToPending. RetryCount matters
+ // as much as status: a peer that took Pending but kept retry_count at max would park
+ // the message again on its first attempt after a failover.
+ await EnqueueAsync(A, "msg-requeue");
+ await SetStatusAsync(A, "msg-requeue", StoreAndForwardMessageStatus.Parked, retryCount: 50);
+ await WaitUntilAsync(
+ () => HasStatusAsync(B, "msg-requeue", StoreAndForwardMessageStatus.Parked),
+ "the message to be parked on B before it is requeued");
+
+ await SetStatusAsync(A, "msg-requeue", StoreAndForwardMessageStatus.Pending, retryCount: 0);
+
+ await WaitUntilAsync(
+ async () => await ReadAsync(B, "msg-requeue")
+ is { Status: (int)StoreAndForwardMessageStatus.Pending, RetryCount: 0 },
+ "node B to show the requeued message as Pending with retry_count 0");
+ }
+
+ [Fact]
+ public async Task MessageAddedThenRemoved_NeverConvergesToPresent()
+ {
+ // The portable intent of ReplicationOperations_AreDispatchedInIssueOrder (see the
+ // class remarks). Ordering only ever mattered because a remove overtaken by its own
+ // add would resurrect a delivered message and send it twice. Under LWW that
+ // reordering is not preventable by dispatch discipline — it is prevented because the
+ // tombstone carries the later HLC and therefore wins regardless of arrival order.
+ //
+ // Asserting the endpoint rather than the sequence is what makes this portable: the
+ // old test would fail on any async transport even when the outcome was correct.
+ await EnqueueAsync(A, "msg-ordering");
+ await DeleteAsync(A, "msg-ordering");
+
+ // A control row written in the same window. Without it this test is VACUOUS: an
+ // absent row is also what a pair that replicates nothing at all looks like, so it
+ // would pass with capture switched off entirely (observed — it was the only one of
+ // these seven that survived unregistering sf_messages). The control converging is
+ // what makes the absence of msg-ordering evidence rather than silence.
+ await EnqueueAsync(A, "msg-ordering-control");
+
+ await WaitUntilAsync(
+ () => ExistsAsync(B, "msg-ordering-control"),
+ "the control message to reach node B, proving the pair was replicating");
+
+ await WaitUntilAsync(
+ async () => !await ExistsAsync(B, "msg-ordering") && !await ExistsAsync(A, "msg-ordering"),
+ "the add-then-remove pair to settle as absent on both nodes");
+
+ // Absence has to persist, not just occur once: a late-arriving add would flip the row
+ // back and the poll above could have observed a gap it never actually converged to.
+ await Task.Delay(TimeSpan.FromSeconds(2));
+ Assert.False(await ExistsAsync(B, "msg-ordering"));
+ Assert.False(await ExistsAsync(A, "msg-ordering"));
+ Assert.True(await ExistsAsync(B, "msg-ordering-control"));
+ }
+
+ [Fact]
+ public async Task SameMessage_BufferedTwice_ConvergesToTheNewerState()
+ {
+ // Was: ApplyReplicatedAdd_Twice_IsIdempotent_NewestWins. Under CDC this is not a
+ // property the application has to implement — LWW on the primary key gives it — but
+ // it is still a property the pair must HAVE, so it is asserted rather than assumed.
+ await EnqueueAsync(A, "msg-twice");
+ await WaitUntilAsync(() => ExistsAsync(B, "msg-twice"), "the first write to reach B");
+
+ await SetStatusAsync(A, "msg-twice", StoreAndForwardMessageStatus.InFlight, retryCount: 3);
+
+ await WaitUntilAsync(
+ async () => await ReadAsync(B, "msg-twice")
+ is { Status: (int)StoreAndForwardMessageStatus.InFlight, RetryCount: 3 },
+ "the newer state to win on node B");
+ }
+
+ [Fact]
+ public async Task ParkArrivingWithoutItsAdd_StillMaterialisesTheRow()
+ {
+ // Was: ApplyReplicatedPark_WhenAddWasLost_MaterializesTheParkedRow. The bespoke
+ // replicator needed explicit upsert semantics so a lost Add did not leave a Park with
+ // nothing to update. The CDC equivalent is a peer that was offline for the add and
+ // only ever sees the row in its parked state — it must still end up with the row,
+ // not skip it as an update to something it never had.
+ await StopPassiveAsync();
+
+ await EnqueueAsync(A, "msg-park-no-add");
+ await SetStatusAsync(A, "msg-park-no-add", StoreAndForwardMessageStatus.Parked, retryCount: 50);
+
+ await RestartPairAsync();
+
+ await WaitUntilAsync(
+ () => HasStatusAsync(B, "msg-park-no-add", StoreAndForwardMessageStatus.Parked),
+ "node B to materialise a row it only ever saw in its parked state");
+ }
+}