diff --git a/src/ZB.MOM.WW.ScadaBridge.Communication/ClusterState/ActiveNodeEvaluator.cs b/src/ZB.MOM.WW.ScadaBridge.Communication/ClusterState/ActiveNodeEvaluator.cs
index 57d7b425..43ba93a7 100644
--- a/src/ZB.MOM.WW.ScadaBridge.Communication/ClusterState/ActiveNodeEvaluator.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.Communication/ClusterState/ActiveNodeEvaluator.cs
@@ -10,11 +10,20 @@ namespace ZB.MOM.WW.ScadaBridge.Communication.ClusterState;
/// once the original first node restarts and rejoins; every product-level active/standby
/// decision must use this evaluator, never cluster.State.Leader.
///
-/// Lives in Communication (not Host) so BOTH SiteCommunicationActor and
-/// SiteReplicationActor can default to it — Host cannot be referenced from either.
-/// The Host's ClusterActivityEvaluator.SelfIsOldest delegates here, so the S&F
-/// delivery gate (IClusterNodeProvider.SelfIsPrimary), the resync authority checks,
-/// and the heartbeat IsActive stamp all share one implementation.
+/// Lives in Communication (not Host) so SiteCommunicationActor can default to it —
+/// Host cannot be referenced from there. The Host's
+/// ClusterActivityEvaluator.SelfIsOldest delegates here, so the S&F delivery gate
+/// (IClusterNodeProvider.SelfIsPrimary) and the heartbeat IsActive stamp share one
+/// implementation.
+///
+/// It also backed SiteReplicationActor's resync authority checks until LocalDb
+/// Phase 2 deleted that actor. Those checks existed because the bespoke resync applied a
+/// destructive delete-all-then-insert-all, so running it in the wrong direction wiped a
+/// live store-and-forward buffer. LocalDb's snapshot resync merges per row under
+/// last-writer-wins and never deletes, so there is no destructive apply left to gate — the
+/// evaluator survives for the delivery gate and the heartbeat, which still genuinely need
+/// a single active node.
+///
///
///
public static class ActiveNodeEvaluator
diff --git a/src/ZB.MOM.WW.ScadaBridge.Host/Actors/AkkaHostedService.cs b/src/ZB.MOM.WW.ScadaBridge.Host/Actors/AkkaHostedService.cs
index 2e29f03a..16b8b41c 100644
--- a/src/ZB.MOM.WW.ScadaBridge.Host/Actors/AkkaHostedService.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.Host/Actors/AkkaHostedService.cs
@@ -766,37 +766,18 @@ akka {{
var deploymentConfigFetcher =
_serviceProvider.GetService();
- // Create SiteReplicationActor on every node (not a singleton)
- var sfStorage = _serviceProvider.GetRequiredService();
- var replicationService = _serviceProvider.GetRequiredService();
- var replicationLogger = _serviceProvider.GetRequiredService()
- .CreateLogger();
-
- // ONE active-node predicate instance governs the S&F delivery gate, the resync
- // authority checks (SiteReplicationActor), and the heartbeat IsActive stamp
- // (SiteCommunicationActor, wired below) — review 02 round 2, N1. Null in
- // non-clustered test hosts: the actors fall back to the shared oldest-Up
- // evaluator, never to a leader check.
+ // ONE active-node predicate instance governs the S&F delivery gate and the
+ // heartbeat IsActive stamp (SiteCommunicationActor, wired below) — review 02
+ // round 2, N1. It also governed SiteReplicationActor's resync authority until
+ // LocalDb Phase 2 deleted that actor: the library's snapshot resync merges per row
+ // under last-writer-wins and never deletes, so there is no destructive apply left
+ // to need an authority check. Null in non-clustered test hosts: the consumers fall
+ // back to the shared oldest-Up evaluator, never to a leader check.
var clusterNodeProvider = _serviceProvider.GetService();
Func? activeNodeCheck = clusterNodeProvider != null
? () => clusterNodeProvider.SelfIsPrimary
: null;
- var replicationActor = _actorSystem!.ActorOf(
- Props.Create(() => new SiteReplicationActor(
- storage, sfStorage, replicationService, siteRole, replicationLogger,
- deploymentConfigFetcher, activeNodeCheck, siteRuntimeOptionsValue, null)),
- "site-replication");
-
- // Wire S&F replication handler to forward operations via the replication actor
- replicationService.SetReplicationHandler(op =>
- {
- replicationActor.Tell(new ReplicateStoreAndForward(op));
- return Task.CompletedTask;
- });
-
- _logger.LogInformation("SiteReplicationActor created and S&F replication handler wired");
-
// Deployment Manager — role-scoped singleton via SingletonRegistrar
// (review 01 round-2 N5): previously hand-rolled with bare PoisonPill
// termination and NO PhaseClusterLeave drain, so in-flight SQLite
@@ -807,7 +788,7 @@ akka {{
_actorSystem!, "deployment-manager",
Props.Create(() => new DeploymentManagerActor(
storage, compilationService, sharedScriptLibrary, streamManager,
- siteRuntimeOptionsValue, dmLogger, dclManager, replicationActor,
+ siteRuntimeOptionsValue, dmLogger, dclManager,
siteHealthCollector, _serviceProvider, null, deploymentConfigFetcher)),
_logger, role: siteRole);
var dmProxy = dm.Proxy;
@@ -1053,7 +1034,7 @@ akka {{
// SetReady asserts a deliberately narrow contract. By this point the
// actor system exists, SiteStreamManager.Initialize has run, and every
// role actor (SiteCommunicationActor, deployment-manager singleton,
- // SiteReplicationActor, the ClusterClient) has been created with ActorOf —
+ // the ClusterClient) has been created with ActorOf —
// creation and the registration Tells are synchronous and strictly ordered.
// What is NOT guaranteed is completion of each actor's PreStart or the
// ClusterClient's initial-contact handshake with central: those are
diff --git a/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbSetup.cs b/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbSetup.cs
index 626c6d5d..beb7f311 100644
--- a/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbSetup.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.Host/SiteLocalDbSetup.cs
@@ -48,8 +48,7 @@ public static class SiteLocalDbSetup
SiteEventLogSchema.Apply(connection);
// Phase 2: the site's configuration tables and the store-and-forward buffer
- // now live in this file too. Created here, deliberately NOT registered — see
- // below.
+ // now live in this file too.
SiteStorageSchema.Apply(connection);
StoreAndForwardSchema.Apply(connection);
}
@@ -60,12 +59,30 @@ public static class SiteLocalDbSetup
db.RegisterReplicated("OperationTracking");
db.RegisterReplicated("site_events");
- // The Phase 2 tables are created above but NOT registered yet. The bespoke
- // SiteReplicationActor / StoreAndForward ReplicationService still own replicating
- // them until the Task 14 cutover deletes both and registers these in one commit.
- // Registering them now would mean two independent replicators writing the same
- // rows — harmless in principle (both upsert) but it would mask a defect in either
- // one, which is exactly what the cutover needs to be able to see.
+ // Phase 2: the store-and-forward buffer and the seven site configuration tables.
+ // These replaced the bespoke SiteReplicationActor and StoreAndForward
+ // ReplicationService, which shipped hand-written Add/Remove/Park/Requeue operations
+ // over Akka; both were deleted in the same commit that added these lines, so the
+ // two mechanisms never ran at once.
+ //
+ // Both composite-PK tables are fine: RegisterReplicated orders multi-column PKs by
+ // ordinal. No Phase 2 table has a BLOB column, which it would reject.
+ db.RegisterReplicated("sf_messages");
+ db.RegisterReplicated("deployed_configurations");
+ db.RegisterReplicated("static_attribute_overrides");
+ db.RegisterReplicated("shared_scripts");
+ db.RegisterReplicated("external_systems");
+ db.RegisterReplicated("database_connections");
+ db.RegisterReplicated("data_connection_definitions");
+ db.RegisterReplicated("native_alarm_state");
+
+ // notification_lists and smtp_configurations are created but deliberately NOT
+ // registered. They are permanently empty by design — the site-side write paths were
+ // removed on 2026-07-10, the legacy migrator skips them, and the active node's
+ // artifact apply purges them on every deploy. Registering them would open a standing
+ // replication channel whose only historical payload was plaintext SMTP passwords, in
+ // exchange for replicating nothing. Anyone adding them here should first establish
+ // that a site has a legitimate reason to hold SMTP credentials at all.
// AFTER registration, so migrated rows enter the oplog and reach the peer like
// any other write. Before it, they would be invisible to replication forever.
diff --git a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Actors/DeploymentManagerActor.cs b/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Actors/DeploymentManagerActor.cs
index e360dffc..56152dc2 100644
--- a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Actors/DeploymentManagerActor.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Actors/DeploymentManagerActor.cs
@@ -52,7 +52,6 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
///
private readonly ILoggerFactory _loggerFactory;
private readonly IActorRef? _dclManager;
- private readonly IActorRef? _replicationActor;
private readonly ISiteHealthCollector? _healthCollector;
private readonly IServiceProvider? _serviceProvider;
///
@@ -166,7 +165,6 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
SiteRuntimeOptions options,
ILogger logger,
IActorRef? dclManager = null,
- IActorRef? replicationActor = null,
ISiteHealthCollector? healthCollector = null,
IServiceProvider? serviceProvider = null,
ILoggerFactory? loggerFactory = null,
@@ -181,7 +179,6 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
_streamManager = streamManager;
_options = options;
_dclManager = dclManager;
- _replicationActor = replicationActor;
_healthCollector = healthCollector;
_serviceProvider = serviceProvider;
_configFetcher = configFetcher;
@@ -785,15 +782,14 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
await _storage.ClearStaticOverridesAsync(instanceName);
await _storage.ClearNativeAlarmsForInstanceAsync(instanceName);
- // Replicate to standby node — notify-and-fetch: send only the deployment id +
- // central fetch coordinates (NOT the config JSON). The standby fetches the
- // config over HTTP itself, so a large config never crosses the intra-site Akka
- // hop (which would silently drop on the 128 KB frame trap). When the coords are
- // absent (deploy paths other than RefreshDeployment), the standby fetch is a
- // no-op miss and reconciliation is the durable backstop.
- _replicationActor?.Tell(new ReplicateConfigDeploy(
- instanceName, command.DeploymentId, command.RevisionHash, true,
- command.CentralFetchBaseUrl ?? "", command.FetchToken ?? ""));
+ // No explicit replication step: deployed_configurations is a replicated table,
+ // so the write above is captured and shipped to the peer like any other row.
+ // This used to be a notify-and-fetch Tell — the standby was sent the deployment
+ // id plus central fetch coordinates and pulled the config over HTTP itself,
+ // because a large config would silently drop on the intra-site Akka hop's
+ // 128 KB frame trap. Replication carries the row directly and has no such limit,
+ // so the standby no longer fetches on deploy. (SiteReconciliationActor still
+ // fetches at node startup when central reports gaps — a different path.)
return new DeployPersistenceResult(
command.DeploymentId, instanceName, true, null, sender, isRedeploy);
@@ -967,7 +963,6 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
{
if (t.IsCompletedSuccessfully)
{
- _replicationActor?.Tell(new ReplicateConfigSetEnabled(instanceName, false));
// Operational `deployment` event — disable succeeded.
LogDeploymentEvent("Info", instanceName, $"Instance {instanceName} disabled");
}
@@ -1001,7 +996,6 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
Task.Run(async () =>
{
await _storage.SetInstanceEnabledAsync(instanceName, true);
- _replicationActor?.Tell(new ReplicateConfigSetEnabled(instanceName, true));
var configs = await _storage.GetAllDeployedConfigsAsync();
var config = configs.FirstOrDefault(c => c.InstanceUniqueName == instanceName);
return new EnableResult(command, config, null, sender);
@@ -1099,7 +1093,6 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
{
if (t.IsCompletedSuccessfully)
{
- _replicationActor?.Tell(new ReplicateConfigRemove(instanceName));
// Operational `deployment` event — delete succeeded.
LogDeploymentEvent("Info", instanceName, $"Instance {instanceName} deleted");
}
@@ -1949,7 +1942,6 @@ public class DeploymentManagerActor : ReceiveActor, IWithTimers
// central-only and is never stored on a site (see the purge above).
// Replicate artifacts to standby node
- _replicationActor?.Tell(new ReplicateArtifacts(command));
return new ArtifactDeploymentResponse(
command.DeploymentId, "", true, null, DateTimeOffset.UtcNow);
diff --git a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Actors/SiteReplicationActor.cs b/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Actors/SiteReplicationActor.cs
deleted file mode 100644
index 63066d2d..00000000
--- a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Actors/SiteReplicationActor.cs
+++ /dev/null
@@ -1,707 +0,0 @@
-using Akka.Actor;
-using Akka.Cluster;
-using Akka.Event;
-using Microsoft.Extensions.Logging;
-using ZB.MOM.WW.ScadaBridge.Commons.Observability;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Deployment;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Messages;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Persistence;
-using ZB.MOM.WW.ScadaBridge.StoreAndForward;
-
-namespace ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors;
-
-///
-/// Runs on every site node (not a singleton). Handles both config and S&F replication
-/// between site cluster peers.
-///
-/// Outbound: receives local replication requests and forwards to peer via ActorSelection.
-/// Inbound: receives replicated operations from peer and applies to local SQLite.
-/// Uses fire-and-forget (Tell) — no ack wait per design.
-///
-public class SiteReplicationActor : ReceiveActor, IWithTimers
-{
- private readonly SiteStorageService _storage;
- private readonly StoreAndForwardStorage _sfStorage;
- private readonly ReplicationService _replicationService;
- private readonly IDeploymentConfigFetcher? _configFetcher;
- private readonly string _siteRole;
- private readonly ILogger _logger;
- private readonly Cluster _cluster;
- private readonly Func _isActive;
- private readonly int _configFetchRetryCount;
- private readonly TimeSpan _configFetchRetryDelay;
- private Address? _peerAddress;
-
- /// Akka timer scheduler injected by the framework via .
- public ITimerScheduler Timers { get; set; } = null!;
-
- // ── Chunked-resync assembly (standby side; actor-thread only) ──
- private string? _assemblingResyncId;
- private int _assemblingTotalChunks;
- private bool _assemblingTruncated;
- private readonly Dictionary> _assemblingChunks = new();
- private const string ResyncAssemblyTimerKey = "sf-resync-assembly-timeout";
-
- /// How long a partial chunk assembly may wait for its missing chunks before
- /// being discarded (a lost chunk = lost resync; the next peer-track retries). Ctor
- /// test seam; production default 30 s.
- private readonly TimeSpan _resyncAssemblyTimeout;
-
- // ── Resync delivery confirmation (active side; actor-thread only) ──
- private string? _pendingAckResyncId;
- private const string ResyncAckTimerKey = "sf-resync-ack-timeout";
-
- /// How long the active node waits for the standby's
- /// before warning + counting the resync as unacknowledged (lost chunks / dead peer). Ctor
- /// test seam; production default 60 s.
- private readonly TimeSpan _resyncAckTimeout;
-
- ///
- /// Maximum rows an active node returns in a single anti-entropy resync snapshot.
- /// A standby whose buffer exceeded this (Truncated snapshot) resyncs the oldest
- /// 10 000 rows; further divergence drains naturally as the active node delivers.
- ///
- private const int MaxResyncRows = 10_000;
-
- ///
- /// Estimated per-chunk payload budget for a resync snapshot. Akka remoting's default
- /// maximum-frame-size is 128 000 bytes and BuildHocon sets no override,
- /// so the monolithic is silently undeliverable for any
- /// realistic backlog (review 02 round 2, N2). 64 000 bytes leaves ≈50% headroom for
- /// the JSON envelope, CLR type manifests, and the non-payload columns.
- ///
- internal const int MaxResyncChunkBytes = 64_000;
-
- /// Row cap per resync chunk (bounds a chunk even when every row is tiny).
- internal const int MaxResyncChunkRows = 200;
-
- ///
- /// Splits a resync snapshot into chunks that fit Akka remoting's default
- /// 128 000-byte frame (review 02 round 2, N2): rows accumulate until the estimated
- /// payload budget or the row cap is hit. Estimation is payload-dominated
- /// (payload_json length + 512 bytes fixed overhead per row); a single row whose
- /// payload exceeds the budget ships alone (Warning at the call site). Order is
- /// preserved (oldest-first, matching GetAllMessagesAsync).
- ///
- internal static List> ChunkForRemoting(
- IReadOnlyList rows, int maxChunkBytes, int maxChunkRows)
- {
- var chunks = new List>();
- var current = new List();
- var currentBytes = 0;
- foreach (var row in rows)
- {
- var estimate = (row.PayloadJson?.Length ?? 0) + 512;
- if (current.Count > 0 && (currentBytes + estimate > maxChunkBytes || current.Count >= maxChunkRows))
- {
- chunks.Add(current);
- current = new List();
- currentBytes = 0;
- }
- current.Add(row);
- currentBytes += estimate;
- }
- if (current.Count > 0) chunks.Add(current);
- return chunks;
- }
-
- ///
- /// Initializes a new and registers Akka message handlers.
- ///
- /// Service for accessing local site storage.
- /// Store-and-forward SQLite storage for replication of buffered messages.
- /// Service providing replication transport logic.
- /// Akka cluster role used to identify peer nodes to replicate to.
- /// Logger instance.
- ///
- /// Fetches a deployed instance's config JSON from central over HTTP. Used by the
- /// notify-and-fetch standby apply path (): the peer
- /// replicates only the deployment id, and the standby fetches the config itself so a large
- /// config never crosses the intra-site Akka hop. Null on nodes/tests without a fetcher.
- ///
- ///
- /// Active-node check that gates the buffer-resync roles (a standby requests a
- /// resync, the active node answers). Production wiring passes the Host's
- /// IClusterNodeProvider.SelfIsPrimary delegate (the same instance gating the
- /// S&F delivery sweep); null falls back to the shared oldest-Up evaluator
- /// ().
- ///
- /// Site runtime options, including the config-fetch retry count; production defaults apply when null.
- /// Delay between config-fetch retry attempts; defaults to 2 seconds when null.
- public SiteReplicationActor(
- SiteStorageService storage,
- StoreAndForwardStorage sfStorage,
- ReplicationService replicationService,
- string siteRole,
- ILogger logger,
- IDeploymentConfigFetcher? configFetcher = null,
- Func? isActiveOverride = null,
- SiteRuntimeOptions? options = null,
- TimeSpan? configFetchRetryDelay = null,
- TimeSpan? resyncAssemblyTimeout = null,
- TimeSpan? resyncAckTimeout = null)
- {
- _storage = storage;
- _sfStorage = sfStorage;
- _replicationService = replicationService;
- _configFetcher = configFetcher;
- _siteRole = siteRole;
- _logger = logger;
- _cluster = Cluster.Get(Context.System);
- _isActive = isActiveOverride ?? DefaultIsActive;
- _resyncAssemblyTimeout = resyncAssemblyTimeout ?? TimeSpan.FromSeconds(30);
- _resyncAckTimeout = resyncAckTimeout ?? TimeSpan.FromSeconds(60);
- // UA2: bound the standby's replicated-config fetch retries. At least one
- // attempt always runs; the fixed inter-attempt delay is a test seam
- // (production default 2 s).
- _configFetchRetryCount = Math.Max(1, options?.ConfigFetchRetryCount ?? 1);
- _configFetchRetryDelay = configFetchRetryDelay ?? TimeSpan.FromSeconds(2);
-
- // Cluster member events
- Receive(HandleMemberUp);
- Receive(HandleMemberRemoved);
- Receive(HandleCurrentClusterState);
-
- // Outbound — forward to peer
- Receive(msg => SendToPeer(new ApplyConfigDeploy(
- msg.InstanceName, msg.DeploymentId, msg.RevisionHash, msg.IsEnabled,
- msg.CentralFetchBaseUrl, msg.FetchToken)));
- Receive(msg => SendToPeer(new ApplyConfigRemove(msg.InstanceName)));
- Receive(msg => SendToPeer(new ApplyConfigSetEnabled(
- msg.InstanceName, msg.IsEnabled)));
- Receive(msg => SendToPeer(new ApplyArtifacts(msg.Command)));
- Receive(msg => SendToPeer(new ApplyStoreAndForward(msg.Operation)));
-
- // Inbound — apply from peer
- Receive(HandleApplyConfigDeploy);
- Receive(HandleApplyConfigRemove);
- Receive(HandleApplyConfigSetEnabled);
- Receive(HandleApplyArtifacts);
- Receive(HandleApplyStoreAndForward);
-
- // Anti-entropy — full S&F buffer resync on peer (re)join
- Receive(HandleRequestSfBufferResync);
- Receive(HandleSfResyncSnapshotLoaded);
- Receive(HandleSfBufferSnapshotChunk);
- Receive(HandleResyncAssemblyTimedOut);
- Receive(msg =>
- {
- if (msg.ResyncId == _pendingAckResyncId)
- {
- _pendingAckResyncId = null;
- Timers.Cancel(ResyncAckTimerKey);
- }
- ScadaBridgeTelemetry.RecordSfResyncCompleted();
- _logger.LogInformation("S&F resync {ResyncId} acknowledged by standby: {Rows} row(s) applied",
- msg.ResyncId, msg.RowCount);
- });
- Receive(msg =>
- {
- if (msg.ResyncId != _pendingAckResyncId) return;
- _pendingAckResyncId = null;
- ScadaBridgeTelemetry.RecordSfResyncAckMissing();
- _logger.LogWarning(
- "S&F resync {ResyncId} was never acknowledged within {Window} — snapshot chunks may have been lost (frame drop / dead peer); the next peer-track retries",
- msg.ResyncId, _resyncAckTimeout);
- });
- Receive(HandleSfBufferSnapshot); // legacy monolithic handler — retained for rolling compat
- }
-
- ///
- protected override void PreStart()
- {
- base.PreStart();
- _cluster.Subscribe(Self, ClusterEvent.SubscriptionInitialStateMode.InitialStateAsSnapshot,
- typeof(ClusterEvent.MemberUp),
- typeof(ClusterEvent.MemberRemoved));
- _logger.LogInformation("SiteReplicationActor started, subscribing to cluster events for role {Role}", _siteRole);
- }
-
- ///
- protected override void PostStop()
- {
- _cluster.Unsubscribe(Self);
- base.PostStop();
- }
-
- private void HandleCurrentClusterState(ClusterEvent.CurrentClusterState state)
- {
- foreach (var member in state.Members)
- {
- if (member.Status == MemberStatus.Up)
- TryTrackPeer(member);
- }
- }
-
- private void HandleMemberUp(ClusterEvent.MemberUp evt)
- {
- TryTrackPeer(evt.Member);
- }
-
- private void HandleMemberRemoved(ClusterEvent.MemberRemoved evt)
- {
- if (evt.Member.Address.Equals(_peerAddress))
- {
- _logger.LogInformation("Peer node removed: {Address}", _peerAddress);
- _peerAddress = null;
- }
- }
-
- private void TryTrackPeer(Member member)
- {
- // Must have our site role, and must not be self
- if (member.HasRole(_siteRole) && !member.Address.Equals(_cluster.SelfAddress))
- {
- _peerAddress = member.Address;
- _logger.LogInformation("Peer node tracked: {Address}", _peerAddress);
- OnPeerTracked();
- }
- }
-
- ///
- /// Side-effect run whenever a peer is (re)tracked. A standby requests a
- /// full S&F buffer snapshot for anti-entropy resync — this closes the "a
- /// standby down for an hour rejoins and diverges forever" gap: it may have missed
- /// replicated Add/Remove/Park ops while it was gone. The active node never
- /// requests. so tests can drive it without a
- /// real two-node cluster.
- ///
- protected virtual void OnPeerTracked()
- {
- if (!SafeIsActive())
- {
- SendToPeer(new RequestSfBufferResync());
- }
- }
-
- ///
- /// Repo-standard active-node check: this node is active when it is the OLDEST Up
- /// member carrying the site role — the same oldest-Up semantics as the S&F delivery
- /// gate (IClusterNodeProvider.SelfIsPrimary → ClusterActivityEvaluator → shared
- /// ActiveNodeEvaluator). NEVER the cluster leader: leadership is lowest-address and
- /// diverges from singleton/delivery placement permanently after the lower-address
- /// node restarts — the divergence that made the delivering node wipe its own live
- /// buffer via a wrong-direction resync (review 02 round 2, N1 Critical). Any other
- /// state reports standby — safe-by-default.
- ///
- private bool DefaultIsActive() =>
- Communication.ClusterState.ActiveNodeEvaluator.SelfIsOldestUp(_cluster, _siteRole);
-
- ///
- /// Evaluates the active-node check, treating a throwing check as standby
- /// (safe-by-default: a standby never delivers or answers resyncs).
- ///
- private bool SafeIsActive()
- {
- try
- {
- return _isActive();
- }
- catch (Exception ex)
- {
- _logger.LogWarning(ex, "Active-node check threw; treating node as standby");
- return false;
- }
- }
-
- ///
- /// Forwards a replication message to the tracked peer node's site-replication actor
- /// (fire-and-forget, dropped when no peer is tracked).
- /// so tests can intercept the peer send without standing up a real two-node cluster.
- ///
- /// The replication message to forward to the peer.
- protected virtual void SendToPeer(object message)
- {
- if (_peerAddress == null)
- {
- // A dropped op is a lost delta — surface it at Warning + a metric so a
- // never-tracked peer (a dead standby) is visible, not silent (arch
- // review 02). In single-node dev the peer is legitimately absent; the
- // per-op warning is rate-tolerable (accepted per review).
- ScadaBridgeTelemetry.RecordReplicationFailure();
- _logger.LogWarning("No peer available, dropping replication message {Type}", message.GetType().Name);
- return;
- }
-
- var path = new RootActorPath(_peerAddress) / "user" / "site-replication";
- Context.ActorSelection(path).Tell(message);
- }
-
- // ── Inbound handlers ──
-
- private void HandleApplyConfigDeploy(ApplyConfigDeploy msg)
- {
- if (string.IsNullOrEmpty(msg.CentralFetchBaseUrl))
- {
- // The direct DeployInstanceCommand cross-cluster wire path was retired.
- // This guard is a defensive fallback: skip quietly rather than calling FetchAsync("")
- // and logging a spurious error. Reconciliation backstops any missed writes.
- _logger.LogDebug(
- "No fetch coords for {Instance} (deployment {DeploymentId}) — skipping replicated fetch; T18 reconciliation is the backstop",
- msg.InstanceName, msg.DeploymentId);
- return;
- }
-
- if (_configFetcher is null)
- {
- _logger.LogWarning(
- "No config fetcher available; cannot apply replicated config for {Instance} (deployment {DeploymentId}) — reconciliation will backstop",
- msg.InstanceName, msg.DeploymentId);
- return;
- }
-
- _logger.LogInformation(
- "Replicating config for {Instance} (deployment {DeploymentId}) — fetching from central",
- msg.InstanceName, msg.DeploymentId);
-
- // Notify-and-fetch: the peer sent only the id, so the standby fetches the config
- // itself (off-thread; best-effort fire-and-forget, matching the no-ack replication
- // model). The guarded write only overwrites a strictly-older local row. The fetch
- // is retried up to ConfigFetchRetryCount times with a fixed delay (UA2) — a transient
- // central hiccup no longer defers to the slower reconciliation backstop, which still
- // covers a total failure after the last attempt.
- _ = FetchWithRetryAsync();
- return;
-
- async Task FetchWithRetryAsync()
- {
- for (var attempt = 1; attempt <= _configFetchRetryCount; attempt++)
- {
- try
- {
- // Non-null: the outer method returns early when _configFetcher is null.
- var json = await _configFetcher!.FetchAsync(
- msg.CentralFetchBaseUrl, msg.DeploymentId, msg.FetchToken, CancellationToken.None);
- await _storage.StoreDeployedConfigIfNewerAsync(
- msg.InstanceName, json, msg.DeploymentId, msg.RevisionHash, msg.IsEnabled);
- return;
- }
- catch (DeploymentConfigFetchException fex) when (fex.IsSuperseded)
- {
- // A superseded/expired fetch never heals by retrying — a newer deploy
- // will replicate its own id.
- _logger.LogInformation(
- "Skip replicated config for {Instance}: superseded/expired (a newer deploy will replicate)",
- msg.InstanceName);
- return;
- }
- catch (Exception ex)
- {
- if (attempt < _configFetchRetryCount)
- {
- _logger.LogWarning(ex,
- "Replicated config fetch attempt {Attempt}/{Max} failed for {Instance} (deployment {DeploymentId}) — retrying in {Delay}",
- attempt, _configFetchRetryCount, msg.InstanceName, msg.DeploymentId, _configFetchRetryDelay);
- await Task.Delay(_configFetchRetryDelay);
- }
- else
- {
- _logger.LogError(ex,
- "Replicated config fetch failed after {Attempts} attempt(s) for {Instance} (deployment {DeploymentId}) — reconciliation will backstop",
- _configFetchRetryCount, msg.InstanceName, msg.DeploymentId);
- }
- }
- }
- }
- }
-
- private void HandleApplyConfigRemove(ApplyConfigRemove msg)
- {
- _logger.LogInformation("Applying replicated config remove for {Instance}", msg.InstanceName);
- _storage.RemoveDeployedConfigAsync(msg.InstanceName)
- .ContinueWith(t =>
- {
- if (t.IsFaulted)
- _logger.LogError(t.Exception, "Failed to apply replicated remove for {Instance}", msg.InstanceName);
- });
- }
-
- private void HandleApplyConfigSetEnabled(ApplyConfigSetEnabled msg)
- {
- _logger.LogInformation("Applying replicated set-enabled={Enabled} for {Instance}", msg.IsEnabled, msg.InstanceName);
- _storage.SetInstanceEnabledAsync(msg.InstanceName, msg.IsEnabled)
- .ContinueWith(t =>
- {
- if (t.IsFaulted)
- _logger.LogError(t.Exception, "Failed to apply replicated set-enabled for {Instance}", msg.InstanceName);
- });
- }
-
- private void HandleApplyArtifacts(ApplyArtifacts msg)
- {
- var command = msg.Command;
- _logger.LogInformation("Applying replicated artifacts, deploymentId={DeploymentId}", command.DeploymentId);
-
- Task.Run(async () =>
- {
- try
- {
- if (command.SharedScripts != null)
- foreach (var s in command.SharedScripts)
- await _storage.StoreSharedScriptAsync(s.Name, s.Code, s.ParameterDefinitions, s.ReturnDefinition);
-
- if (command.ExternalSystems != null)
- foreach (var es in command.ExternalSystems)
- await _storage.StoreExternalSystemAsync(es.Name, es.EndpointUrl, es.AuthType, es.AuthConfiguration, es.MethodDefinitionsJson, es.TimeoutSeconds);
-
- if (command.DatabaseConnections != null)
- foreach (var db in command.DatabaseConnections)
- await _storage.StoreDatabaseConnectionAsync(db.Name, db.ConnectionString, db.MaxRetries, db.RetryDelay);
-
- // Notification lists and SMTP
- // configuration are central-only and are never persisted on a site.
- // Mirror the primary apply path: purge any pre-fix rows (including the
- // plaintext SMTP password) instead of writing the command's
- // (now-always-null) NotificationLists/SmtpConfigurations.
- await _storage.PurgeCentralOnlyNotificationConfigAsync();
-
- if (command.DataConnections != null)
- foreach (var dc in command.DataConnections)
- await _storage.StoreDataConnectionDefinitionAsync(dc.Name, dc.Protocol, dc.PrimaryConfigurationJson, dc.BackupConfigurationJson, dc.FailoverRetryCount);
- }
- catch (Exception ex)
- {
- _logger.LogError(ex, "Failed to apply replicated artifacts");
- }
- });
- }
-
- private void HandleApplyStoreAndForward(ApplyStoreAndForward msg)
- {
- _logger.LogDebug("Applying replicated S&F operation {OpType} for message {Id}",
- msg.Operation.OperationType, msg.Operation.MessageId);
-
- _replicationService.ApplyReplicatedOperationAsync(msg.Operation, _sfStorage)
- .ContinueWith(t =>
- {
- if (t.IsFaulted)
- _logger.LogError(t.Exception, "Failed to apply replicated S&F operation {Id}", msg.Operation.MessageId);
- });
- }
-
- ///
- /// Active-node side of the anti-entropy resync: answers a standby's
- /// with a sequence of byte-budgeted
- /// s (up to oldest rows).
- /// A non-active node ignores the request — only the authoritative node may answer.
- /// The snapshot is piped back to Self so chunking + ack bookkeeping stays actor-safe.
- ///
- private void HandleRequestSfBufferResync(RequestSfBufferResync msg)
- {
- if (!SafeIsActive())
- {
- _logger.LogDebug("Ignoring S&F buffer resync request — this node is not active");
- return;
- }
-
- var replyTo = Sender;
- _sfStorage.GetAllMessagesAsync(MaxResyncRows).PipeTo(
- Self,
- failure: ex => new Status.Failure(ex),
- success: result => new SfResyncSnapshotLoaded(replyTo, result.Messages, result.Truncated));
- }
-
- ///
- /// Active-node continuation: the resync snapshot finished loading; chunk it to fit the
- /// remoting frame and send the sequenced chunks to the requester (all sharing one
- /// resyncId). Task 7 arms the ack-timeout here.
- ///
- private void HandleSfResyncSnapshotLoaded(SfResyncSnapshotLoaded msg)
- {
- var chunks = ChunkForRemoting(msg.Messages, MaxResyncChunkBytes, MaxResyncChunkRows);
- if (chunks.Count == 0) chunks.Add(new List()); // empty buffer still resyncs (clears the standby)
- var resyncId = Guid.NewGuid().ToString("N");
- for (var i = 0; i < chunks.Count; i++)
- {
- if (chunks[i].Count == 1 && (chunks[i][0].PayloadJson?.Length ?? 0) + 512 > MaxResyncChunkBytes)
- _logger.LogWarning(
- "Resync row {Id} alone exceeds the chunk budget ({Bytes}B payload); sending solo — it may exceed the remoting frame",
- chunks[i][0].Id, chunks[i][0].PayloadJson?.Length ?? 0);
- msg.ReplyTo.Tell(new SfBufferSnapshotChunk(resyncId, i + 1, chunks.Count, chunks[i], msg.Truncated), Self);
- }
- _logger.LogInformation(
- "Answered S&F resync request with {Rows} row(s) in {Chunks} chunk(s), resyncId={ResyncId}",
- msg.Messages.Count, chunks.Count, resyncId);
- // Arm the ack window: absence of an SfBufferResyncAck within it surfaces as a
- // Warning + counter (the silent-loss mode N2 flagged). Single-outstanding-resync
- // bookkeeping: a new request supersedes by overwriting the id and restarting the timer.
- _pendingAckResyncId = resyncId;
- Timers.StartSingleTimer(ResyncAckTimerKey, new ResyncAckTimedOut(resyncId), _resyncAckTimeout);
- }
-
- ///
- /// Standby-node side of the anti-entropy resync: replaces the local buffer
- /// wholesale with the active node's snapshot. Combined with the upsert-based
- /// replicated applies (arch review 02), any replicated op that lands
- /// after this resync merges cleanly onto the resynced state. An active node
- /// ignores a snapshot — it is the source of truth, never a resync target.
- ///
- private void HandleSfBufferSnapshot(SfBufferSnapshot msg)
- {
- if (SafeIsActive())
- {
- _logger.LogDebug("Ignoring S&F buffer snapshot — this node is active");
- return;
- }
-
- if (msg.Truncated)
- {
- _logger.LogWarning(
- "S&F buffer resync snapshot truncated at {Cap} rows; divergence beyond the cap drains naturally",
- MaxResyncRows);
- }
-
- _logger.LogInformation(
- "Applying S&F buffer resync snapshot ({Count} rows), replacing local buffer", msg.Messages.Count);
-
- Task.Run(async () =>
- {
- // Belt-and-braces (N1): re-check at apply time. ReplaceAllAsync discards
- // every in-flight row (StoreAndForwardStorage.cs "Never call on an active
- // node"); a flip between message receipt and this point must abort.
- if (SafeIsActive())
- {
- _logger.LogWarning(
- "Discarding S&F buffer resync snapshot: this node became active before apply");
- return;
- }
- await _sfStorage.ReplaceAllAsync(msg.Messages);
- })
- .ContinueWith(t =>
- {
- if (t.IsFaulted)
- _logger.LogError(t.Exception, "Failed to apply S&F buffer resync snapshot");
- });
- }
-
- ///
- /// Standby-node side of the chunked anti-entropy resync: accumulates the chunks of one
- /// ResyncId, and once all have arrived, assembles them in sequence order and
- /// replaces the local buffer atomically, then acks. A new ResyncId discards any
- /// stale partial assembly (review 02 round 2, N2). An active node ignores chunks.
- ///
- private void HandleSfBufferSnapshotChunk(SfBufferSnapshotChunk msg)
- {
- if (SafeIsActive())
- {
- _logger.LogDebug("Ignoring S&F resync chunk — this node is active");
- return;
- }
-
- if (_assemblingResyncId != msg.ResyncId)
- {
- if (_assemblingResyncId != null)
- _logger.LogWarning(
- "Discarding partial S&F resync assembly {Old} ({Have}/{Want} chunks): a new resync {New} superseded it",
- _assemblingResyncId, _assemblingChunks.Count, _assemblingTotalChunks, msg.ResyncId);
- _assemblingResyncId = msg.ResyncId;
- _assemblingTotalChunks = msg.TotalChunks;
- _assemblingTruncated = msg.Truncated;
- _assemblingChunks.Clear();
- }
-
- _assemblingChunks[msg.Sequence] = msg.Messages;
- Timers.StartSingleTimer(ResyncAssemblyTimerKey, new ResyncAssemblyTimedOut(msg.ResyncId), _resyncAssemblyTimeout);
-
- if (_assemblingChunks.Count < _assemblingTotalChunks)
- return;
-
- // Complete: assemble in sequence order and apply atomically.
- var assembled = Enumerable.Range(1, _assemblingTotalChunks)
- .SelectMany(seq => _assemblingChunks[seq])
- .ToList();
- var resyncId = _assemblingResyncId!;
- var truncated = _assemblingTruncated;
- _assemblingResyncId = null;
- _assemblingChunks.Clear();
- Timers.Cancel(ResyncAssemblyTimerKey);
-
- if (truncated)
- _logger.LogWarning(
- "S&F buffer resync snapshot truncated at {Cap} rows; divergence beyond the cap drains naturally",
- MaxResyncRows);
- _logger.LogInformation(
- "Applying chunked S&F resync {ResyncId} ({Count} rows), replacing local buffer", resyncId, assembled.Count);
-
- // KNOWN, ACCEPTED race (review 02 round 2, N5 — do NOT "fix" this into something
- // worse): a replicated Remove sent after the active node read its snapshot but
- // before the snapshot's chunks is ordered BEFORE them on the wire (same
- // sender/receiver pair), so this apply can re-add the removed row → an orphan
- // Pending row that a later failover re-delivers ONCE. Bounded, self-correcting at
- // the next resync, and inherent to no-ack replication; a delivered-side dedup or
- // op-sequencing scheme would cost far more than one rare duplicate.
- var replyTo = Sender;
- Task.Run(async () =>
- {
- if (SafeIsActive()) // belt-and-braces, mirrors the monolithic path (T3)
- {
- _logger.LogWarning("Discarding chunked S&F resync {ResyncId}: node became active before apply", resyncId);
- return;
- }
- await _sfStorage.ReplaceAllAsync(assembled);
- replyTo.Tell(new SfBufferResyncAck(resyncId, assembled.Count));
- })
- .ContinueWith(t =>
- {
- if (t.IsFaulted)
- _logger.LogError(t.Exception, "Failed to apply chunked S&F resync {ResyncId}", resyncId);
- });
- }
-
- private void HandleResyncAssemblyTimedOut(ResyncAssemblyTimedOut msg)
- {
- if (_assemblingResyncId != msg.ResyncId) return; // superseded already
- _logger.LogWarning(
- "S&F resync assembly {ResyncId} timed out with {Have}/{Want} chunks — discarding partial (a lost chunk; the next peer-track retries)",
- msg.ResyncId, _assemblingChunks.Count, _assemblingTotalChunks);
- ScadaBridgeTelemetry.RecordReplicationFailure();
- _assemblingResyncId = null;
- _assemblingChunks.Clear();
- }
-
- /// Internal: the resync snapshot finished loading; chunk and send to the requester.
- internal sealed record SfResyncSnapshotLoaded(
- IActorRef ReplyTo, List Messages, bool Truncated);
-
- /// Internal: a partial chunk assembly for exceeded its window.
- internal sealed record ResyncAssemblyTimedOut(string ResyncId);
-
- /// Internal: the active node's ack window for expired.
- internal sealed record ResyncAckTimedOut(string ResyncId);
-}
-
-///
-/// Standby→active: request a full S&F buffer snapshot for anti-entropy resync
-/// (sent when a standby (re)tracks a peer). Crosses Akka remoting between the two
-/// site nodes; the POCO rides the default serializer.
-///
-public sealed record RequestSfBufferResync;
-
-///
-/// Active→standby: full-buffer snapshot. is true when
-/// the active node's buffer exceeded MaxResyncRows (the standby logs a Warning —
-/// divergence beyond the cap drains naturally as the active node delivers). Crosses
-/// Akka remoting; the message list rides the default serializer.
-///
-public sealed record SfBufferSnapshot(List Messages, bool Truncated);
-
-///
-/// Active→standby: one sequenced chunk of a full-buffer anti-entropy snapshot
-/// (review 02 round 2, N2 — the monolithic exceeds Akka
-/// remoting's default 128 000-byte frame for any realistic backlog). All chunks of one
-/// resync share ; is 1-based up to
-/// . Additive message — the legacy monolithic snapshot
-/// handler is retained for rolling upgrades. Crosses intra-site Akka remoting (NOT
-/// ClusterClient — ClusterClientContractLockTests is intentionally not involved).
-///
-public sealed record SfBufferSnapshotChunk(
- string ResyncId, int Sequence, int TotalChunks,
- List Messages, bool Truncated);
-
-///
-/// Standby→active: delivery confirmation — the standby assembled all chunks of
-/// and applied them atomically (
-/// rows installed). Absence within the ack window is surfaced by the active node
-/// (Warning + counter) — the silent-loss mode N2 flagged.
-///
-public sealed record SfBufferResyncAck(string ResyncId, int RowCount);
diff --git a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Messages/ReplicationMessages.cs b/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Messages/ReplicationMessages.cs
deleted file mode 100644
index b937cc0d..00000000
--- a/src/ZB.MOM.WW.ScadaBridge.SiteRuntime/Messages/ReplicationMessages.cs
+++ /dev/null
@@ -1,44 +0,0 @@
-using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
-using ZB.MOM.WW.ScadaBridge.StoreAndForward;
-
-namespace ZB.MOM.WW.ScadaBridge.SiteRuntime.Messages;
-
-// Outbound messages — sent by local DeploymentManagerActor/S&F service
-// to the local SiteReplicationActor for forwarding to the peer node.
-
-/// Outbound: tell the peer to fetch+apply a deployed instance config by id (notify-and-fetch; no inline config).
-public record ReplicateConfigDeploy(
- string InstanceName, string DeploymentId, string RevisionHash, bool IsEnabled,
- string CentralFetchBaseUrl, string FetchToken);
-
-/// Outbound: replicate removal of a deployed instance config to the peer node.
-public record ReplicateConfigRemove(string InstanceName);
-
-/// Outbound: replicate an instance enabled/disabled flag change to the peer node.
-public record ReplicateConfigSetEnabled(string InstanceName, bool IsEnabled);
-
-/// Outbound: replicate a system-wide artifact deployment (shared scripts, external systems, etc.) to the peer node.
-public record ReplicateArtifacts(DeployArtifactsCommand Command);
-
-/// Outbound: replicate a store-and-forward buffer mutation (enqueue/dequeue/park/etc.) to the peer node.
-public record ReplicateStoreAndForward(ReplicationOperation Operation);
-
-// Inbound messages — received from the peer's SiteReplicationActor
-// and applied to local SQLite storage.
-
-/// Inbound: peer-replicated config deploy — the standby fetches the config by id and writes it (guarded).
-public record ApplyConfigDeploy(
- string InstanceName, string DeploymentId, string RevisionHash, bool IsEnabled,
- string CentralFetchBaseUrl, string FetchToken);
-
-/// Inbound: apply peer-replicated removal of a deployed instance config to local SQLite.
-public record ApplyConfigRemove(string InstanceName);
-
-/// Inbound: apply a peer-replicated instance enabled/disabled flag change to local SQLite.
-public record ApplyConfigSetEnabled(string InstanceName, bool IsEnabled);
-
-/// Inbound: apply a peer-replicated system-wide artifact deployment to local SQLite.
-public record ApplyArtifacts(DeployArtifactsCommand Command);
-
-/// Inbound: apply a peer-replicated store-and-forward buffer mutation to the local buffer.
-public record ApplyStoreAndForward(ReplicationOperation Operation);
diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/ReplicationService.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/ReplicationService.cs
deleted file mode 100644
index 60f2b162..00000000
--- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/ReplicationService.cs
+++ /dev/null
@@ -1,194 +0,0 @@
-using Microsoft.Extensions.Logging;
-using ZB.MOM.WW.ScadaBridge.Commons.Observability;
-using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
-
-namespace ZB.MOM.WW.ScadaBridge.StoreAndForward;
-
-///
-/// Async replication of buffer operations to standby node.
-///
-/// - Forwards add/remove/park operations to standby via a replication handler.
-/// - No ack wait (fire-and-forget per design).
-/// - Standby applies operations to its own SQLite.
-/// - On failover, standby resumes delivery from its replicated state.
-///
-public class ReplicationService
-{
- private readonly StoreAndForwardOptions _options;
- private readonly ILogger _logger;
- private Func? _replicationHandler;
-
- /// Initializes a new instance of .
- /// Store-and-forward configuration options.
- /// Logger instance.
- public ReplicationService(
- StoreAndForwardOptions options,
- ILogger logger)
- {
- _options = options;
- _logger = logger;
- }
-
- ///
- /// Sets the handler for forwarding replication operations to the standby node.
- /// Typically wraps Akka Tell to the standby's replication actor.
- ///
- /// The async delegate that forwards each replication operation to the standby.
- public void SetReplicationHandler(Func handler)
- {
- _replicationHandler = handler;
- }
-
- ///
- /// Replicates an enqueue operation to standby (fire-and-forget).
- ///
- /// The message that was enqueued on the active node.
- public void ReplicateEnqueue(StoreAndForwardMessage message)
- {
- if (!_options.ReplicationEnabled || _replicationHandler == null) return;
-
- FireAndForget(new ReplicationOperation(
- ReplicationOperationType.Add,
- message.Id,
- message));
- }
-
- ///
- /// Replicates a remove operation to standby (fire-and-forget).
- ///
- /// The identifier of the message to remove from the standby buffer.
- public void ReplicateRemove(string messageId)
- {
- if (!_options.ReplicationEnabled || _replicationHandler == null) return;
-
- FireAndForget(new ReplicationOperation(
- ReplicationOperationType.Remove,
- messageId,
- null));
- }
-
- ///
- /// Replicates a park operation to standby (fire-and-forget).
- ///
- /// The message that was parked on the active node.
- public void ReplicatePark(StoreAndForwardMessage message)
- {
- if (!_options.ReplicationEnabled || _replicationHandler == null) return;
-
- FireAndForget(new ReplicationOperation(
- ReplicationOperationType.Park,
- message.Id,
- message));
- }
-
- ///
- /// Replicates an operator-initiated requeue (a parked
- /// message moved back to the pending queue) to standby (fire-and-forget). The
- /// carried message reflects the active node's post-requeue state (Pending,
- /// retry_count = 0) so the standby's copy can be brought into sync.
- ///
- /// The message in its post-requeue (Pending, retry_count=0) state.
- public void ReplicateRequeue(StoreAndForwardMessage message)
- {
- if (!_options.ReplicationEnabled || _replicationHandler == null) return;
-
- FireAndForget(new ReplicationOperation(
- ReplicationOperationType.Requeue,
- message.Id,
- message));
- }
-
- ///
- /// Applies a replicated operation received from the active node.
- /// Used by the standby node to keep its SQLite in sync.
- ///
- /// Add/Park/Requeue are applied as upserts (),
- /// not blind INSERT/UPDATE: the full message rides in every one of those operations,
- /// so a Park/Requeue whose original Add was lost (fire-and-forget replication is
- /// best-effort) self-heals by materialising the row, and a duplicate Add (e.g.
- /// re-issued by an anti-entropy resync) applies newest-wins instead of throwing a
- /// primary-key violation. Remove is a plain delete.
- ///
- /// The replication operation to apply.
- /// The standby node's store-and-forward storage to update.
- /// A task representing the asynchronous apply operation.
- public async Task ApplyReplicatedOperationAsync(
- ReplicationOperation operation,
- StoreAndForwardStorage storage)
- {
- switch (operation.OperationType)
- {
- case ReplicationOperationType.Add when operation.Message != null:
- await storage.UpsertMessageAsync(operation.Message);
- break;
-
- case ReplicationOperationType.Remove:
- await storage.RemoveMessageAsync(operation.MessageId);
- break;
-
- case ReplicationOperationType.Park when operation.Message != null:
- operation.Message.Status = StoreAndForwardMessageStatus.Parked;
- await storage.UpsertMessageAsync(operation.Message);
- break;
-
- case ReplicationOperationType.Requeue when operation.Message != null:
- operation.Message.Status = StoreAndForwardMessageStatus.Pending;
- operation.Message.RetryCount = 0;
- await storage.UpsertMessageAsync(operation.Message);
- break;
- }
- }
-
- private void FireAndForget(ReplicationOperation operation)
- {
- // Invoked inline, NOT via Task.Run: the handler is a non-blocking Akka
- // Tell, and thread-pool hand-off destroyed Add/Remove ordering for the
- // same message id (arch review 02). Inline invocation preserves issue
- // order; Akka's per-sender/receiver guarantee preserves it on the wire.
- try
- {
- var task = _replicationHandler!.Invoke(operation);
- if (!task.IsCompletedSuccessfully)
- {
- task.ContinueWith(t =>
- {
- ScadaBridgeTelemetry.RecordReplicationFailure();
- _logger.LogWarning(t.Exception,
- "Replication of {OpType} for message {MessageId} failed (best-effort); standby buffer may be diverging",
- operation.OperationType, operation.MessageId);
- },
- TaskContinuationOptions.OnlyOnFaulted);
- }
- }
- catch (Exception ex)
- {
- ScadaBridgeTelemetry.RecordReplicationFailure();
- _logger.LogWarning(ex,
- "Replication of {OpType} for message {MessageId} failed (best-effort); standby buffer may be diverging",
- operation.OperationType, operation.MessageId);
- }
- }
-}
-
-///
-/// Represents a buffer operation to be replicated to standby.
-///
-public record ReplicationOperation(
- ReplicationOperationType OperationType,
- string MessageId,
- StoreAndForwardMessage? Message);
-
-///
-/// Types of buffer operations that are replicated.
-///
-public enum ReplicationOperationType
-{
- Add,
- Remove,
- Park,
- ///
- /// An operator moved a parked message back to the pending
- /// queue. The standby resets its matching row to Pending with retry_count = 0.
- ///
- Requeue
-}
diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/ServiceCollectionExtensions.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/ServiceCollectionExtensions.cs
index a5c2e513..314b5522 100644
--- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/ServiceCollectionExtensions.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/ServiceCollectionExtensions.cs
@@ -28,7 +28,6 @@ public static class ServiceCollectionExtensions
var storage = sp.GetRequiredService();
var options = sp.GetRequiredService>().Value;
var logger = sp.GetRequiredService>();
- var replication = sp.GetRequiredService();
// Wire the cached-call lifecycle
// observer + site identity through DI so the S&F retry loop emits
// per-attempt + terminal telemetry under the same TrackedOperationId
@@ -53,19 +52,11 @@ public static class ServiceCollectionExtensions
storage,
options,
logger,
- replication,
cachedCallObserver,
siteId,
siteEventLogger);
});
- services.AddSingleton(sp =>
- {
- var options = sp.GetRequiredService>().Value;
- var logger = sp.GetRequiredService>();
- return new ReplicationService(options, logger);
- });
-
return services;
}
diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs
index 2b8dda17..53a6ad2f 100644
--- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs
@@ -36,7 +36,6 @@ public class StoreAndForwardService
{
private readonly StoreAndForwardStorage _storage;
private readonly StoreAndForwardOptions _options;
- private readonly ReplicationService? _replication;
private readonly ILogger _logger;
///
/// Site-side observer notified
@@ -107,7 +106,7 @@ public class StoreAndForwardService
/// null when no sweep is currently running. Captured when the timer
/// callback starts a sweep so can wait for it to
/// finish before the host disposes downstream dependencies
- /// (, ) that the sweep is
+ /// () that the sweep is
/// still touching. Written from the timer thread and from
/// , so reads are synchronised via the
/// APIs.
@@ -227,7 +226,6 @@ public class StoreAndForwardService
/// The storage backend for buffered messages.
/// Configuration options.
/// Logger instance.
- /// Optional replication service for standby synchronization.
/// Optional observer for cached call lifecycle events.
/// The site identifier this service belongs to.
///
@@ -240,7 +238,6 @@ public class StoreAndForwardService
StoreAndForwardStorage storage,
StoreAndForwardOptions options,
ILogger logger,
- ReplicationService? replication = null,
ICachedCallLifecycleObserver? cachedCallObserver = null,
string siteId = "",
ISiteEventLogger? siteEventLogger = null)
@@ -248,7 +245,6 @@ public class StoreAndForwardService
_storage = storage;
_options = options;
_logger = logger;
- _replication = replication;
_cachedCallObserver = cachedCallObserver;
_siteId = string.IsNullOrWhiteSpace(siteId) ? UnknownSiteSentinel : siteId;
_siteEventLogger = siteEventLogger;
@@ -651,7 +647,6 @@ public class StoreAndForwardService
private async Task BufferAsync(StoreAndForwardMessage message)
{
await _storage.EnqueueAsync(message);
- _replication?.ReplicateEnqueue(message);
Interlocked.Increment(ref _bufferedCount);
}
@@ -802,7 +797,6 @@ public class StoreAndForwardService
if (success)
{
await _storage.RemoveMessageAsync(message.Id);
- _replication?.ReplicateRemove(message.Id);
Interlocked.Decrement(ref _bufferedCount);
RaiseActivity("Delivered", message.Category,
$"Delivered to {message.Target} after {message.RetryCount} retries");
@@ -833,7 +827,6 @@ public class StoreAndForwardService
return RetryOutcome.Skipped;
}
Interlocked.Decrement(ref _bufferedCount);
- _replication?.ReplicatePark(message);
RaiseActivity("Parked", message.Category,
$"Permanent failure for {message.Target}: handler returned false");
@@ -869,7 +862,6 @@ public class StoreAndForwardService
return RetryOutcome.Skipped;
}
Interlocked.Decrement(ref _bufferedCount);
- _replication?.ReplicatePark(message);
RaiseActivity("Parked", message.Category,
$"Max retries ({message.MaxRetries}) reached for {message.Target}");
_logger.LogWarning(
@@ -1077,17 +1069,12 @@ public class StoreAndForwardService
///
/// Retries a parked message (moves back to pending queue).
///
- /// An operator requeue is a buffer state change and is
- /// replicated to the standby (as a )
- /// so a failover preserves the operator's retry intent.
+ /// An operator requeue is a buffer state change, and the peer picks it up because
+ /// sf_messages is a replicated table: the row update is captured and shipped
+ /// like any other write, so a failover preserves the operator's retry intent. This
+ /// used to be an explicit Requeue operation sent to the standby.
/// The activity-log entry carries the message's true
/// category rather than a hard-coded one.
- /// The parked row is captured before the local
- /// requeue write rather than re-read after it, so a concurrent
- /// RemoveMessageAsync or DiscardParkedMessageAsync running
- /// between the two storage calls cannot leave the standby in Parked
- /// while the active node has already requeued — we always have the row in
- /// hand for the Requeue replication.
///
/// The identifier of the message to retry.
/// True if successfully retried, false otherwise.
@@ -1117,7 +1104,6 @@ public class StoreAndForwardService
captured.RetryCount = 0;
captured.LastError = null;
captured.LastAttemptAt = null;
- _replication?.ReplicateRequeue(captured);
RaiseActivity("Retry", captured.Category,
$"Parked message {messageId} moved back to queue");
@@ -1127,9 +1113,10 @@ public class StoreAndForwardService
///
/// Permanently discards a parked message.
///
- /// An operator discard is a buffer removal and is replicated
- /// to the standby (as a ) so the
- /// discarded message does not reappear after a failover.
+ /// An operator discard is a buffer removal. The delete is captured as a tombstone on
+ /// the replicated sf_messages table and carries the later HLC, so the discarded
+ /// message cannot reappear after a failover regardless of arrival order. This used to
+ /// be an explicit Remove operation sent to the standby.
/// The activity-log entry carries the message's true
/// category rather than a hard-coded one.
///
@@ -1143,7 +1130,6 @@ public class StoreAndForwardService
var success = await _storage.DiscardParkedMessageAsync(messageId);
if (success)
{
- _replication?.ReplicateRemove(messageId);
RaiseActivity("Discard", message?.Category ?? StoreAndForwardCategory.ExternalSystem,
$"Parked message {messageId} discarded");
}
diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardStorage.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardStorage.cs
index fbe87cc6..54859796 100644
--- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardStorage.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardStorage.cs
@@ -65,7 +65,7 @@ public class StoreAndForwardStorage
///
/// INSERT statement for a full message row. Shared by
- /// and ; bind with
+ /// and ; bind with
/// so the column list and the parameters never drift apart.
///
private const string InsertMessageSql = @"
@@ -151,39 +151,6 @@ public class StoreAndForwardStorage
return (rows.Take(limit).ToList(), truncated);
}
- ///
- /// Atomically replaces the entire buffer with in a
- /// single transaction (delete-all then insert-all). Standby-side anti-entropy
- /// apply: a peer-join resync overwrites the standby's divergent copy with the
- /// active node's authoritative snapshot. Never call on an active node —
- /// it discards every in-flight row.
- ///
- /// The full buffer snapshot to install.
- /// A task that represents the asynchronous operation.
- public async Task ReplaceAllAsync(IReadOnlyList messages)
- {
- await using var connection = OpenConnection();
- await using var transaction = (SqliteTransaction)await connection.BeginTransactionAsync();
-
- await using (var deleteCmd = connection.CreateCommand())
- {
- deleteCmd.Transaction = transaction;
- deleteCmd.CommandText = "DELETE FROM sf_messages";
- await deleteCmd.ExecuteNonQueryAsync();
- }
-
- foreach (var message in messages)
- {
- await using var insertCmd = connection.CreateCommand();
- insertCmd.Transaction = transaction;
- insertCmd.CommandText = InsertMessageSql;
- BindMessageParameters(insertCmd, message);
- await insertCmd.ExecuteNonQueryAsync();
- }
-
- await transaction.CommitAsync();
- }
-
///
/// Inserts a message, or updates every mutable column in place if a row with the
/// same id already exists (ON CONFLICT(id) DO UPDATE). Used by the standby
diff --git a/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/CompositionRootTests.cs b/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/CompositionRootTests.cs
index 488f186f..59d4bf31 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/CompositionRootTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/CompositionRootTests.cs
@@ -471,7 +471,6 @@ public class SiteCompositionRootTests : IDisposable
new object[] { typeof(IDataConnectionFactory) },
new object[] { typeof(StoreAndForwardStorage) },
new object[] { typeof(StoreAndForwardService) },
- new object[] { typeof(ReplicationService) },
new object[] { typeof(ISiteEventLogger) },
new object[] { typeof(IEventLogQueryService) },
new object[] { typeof(ISiteIdentityProvider) },
diff --git a/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbWiringTests.cs b/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbWiringTests.cs
index f9e8e3eb..4340a148 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbWiringTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.Host.Tests/SiteLocalDbWiringTests.cs
@@ -243,22 +243,65 @@ public class SiteLocalDbWiringTests : IDisposable
}
[Fact]
- public void Site_LocalDb_DoesNotYetRegisterThePhase2Tables()
+ public void Site_LocalDb_RegistersExactlyTheTenReplicatedTables()
{
- // The "not yet" is the assertion that matters, and it is why this is an EXACTLY
- // check rather than a Contains. Until the Task 14 cutover, the bespoke
- // SiteReplicationActor and the StoreAndForward ReplicationService still own these
- // tables. Registering them early would run two replicators over the same rows and
- // hide a defect in either one behind the other's writes.
+ // The Phase 2 cutover. This is an EXACTLY check rather than a Contains in both
+ // directions, and both directions are load-bearing.
//
- // When Task 14 lands, this test does not get deleted — it gets inverted.
+ // Too few: the table silently stops replicating. Storage still works, every other
+ // test still passes, and the pair just quietly diverges — the failure mode Phase 1
+ // existed to end.
+ //
+ // Too many: notification_lists and smtp_configurations must NOT be here. They are
+ // permanently empty by design, and registering them would open a standing
+ // replication channel whose only historical payload was plaintext SMTP passwords.
+ // A Contains-based test would never catch that.
var db = _host.Services.GetRequiredService();
Assert.Equal(
- ["OperationTracking", "site_events"],
+ [
+ // Ordinal order: '_' (0x5F) sorts before 'b' (0x62), so
+ // data_connection_definitions precedes database_connections.
+ "OperationTracking", "data_connection_definitions", "database_connections",
+ "deployed_configurations", "external_systems", "native_alarm_state",
+ "sf_messages", "shared_scripts", "site_events", "static_attribute_overrides",
+ ],
db.ReplicatedTables.Keys.OrderBy(k => k, StringComparer.Ordinal).ToArray());
}
+ [Fact]
+ public void Site_LocalDb_DoesNotReplicateTheCentralOnlyNotificationTables()
+ {
+ // Stated separately from the exact-set test above because this one is a security
+ // property, not a wiring property, and deserves to fail with its own name. The
+ // tables exist (SiteStorageSchema creates them) — they simply must never be
+ // captured. See SiteLocalDbSetup.OnReady for the rationale.
+ var db = _host.Services.GetRequiredService();
+
+ Assert.DoesNotContain("notification_lists", db.ReplicatedTables.Keys);
+ Assert.DoesNotContain("smtp_configurations", db.ReplicatedTables.Keys);
+ }
+
+ [Fact]
+ public void Site_LocalDb_ReplicatedPhase2TablesHaveTheirExpectedPrimaryKeys()
+ {
+ // RegisterReplicated refuses a table with no explicit PK, so reaching these
+ // assertions proves the DDL ran before registration. The composite keys are the
+ // interesting ones: LWW conflict resolution keys on the FULL PK, so a wrong or
+ // truncated key set would silently collapse distinct rows into one.
+ var db = _host.Services.GetRequiredService();
+
+ Assert.Equal(["id"], db.ReplicatedTables["sf_messages"].PkColumns);
+ Assert.Equal(
+ ["instance_unique_name"], db.ReplicatedTables["deployed_configurations"].PkColumns);
+ Assert.Equal(
+ ["instance_unique_name", "attribute_name"],
+ db.ReplicatedTables["static_attribute_overrides"].PkColumns);
+ Assert.Equal(
+ ["instance_unique_name", "source_canonical_name", "source_reference"],
+ db.ReplicatedTables["native_alarm_state"].PkColumns);
+ }
+
private static HashSet TableNames(ILocalDb db)
{
using var connection = db.CreateConnection();
diff --git a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/Cluster/SfBufferResyncPredicateTests.cs b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/Cluster/SfBufferResyncPredicateTests.cs
deleted file mode 100644
index 9e9f3cb2..00000000
--- a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/Cluster/SfBufferResyncPredicateTests.cs
+++ /dev/null
@@ -1,122 +0,0 @@
-using Akka.Actor;
-using Microsoft.Extensions.Logging.Abstractions;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Persistence;
-using ZB.MOM.WW.ScadaBridge.StoreAndForward;
-using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
-using ZB.MOM.WW.ScadaBridge.TestSupport;
-
-namespace ZB.MOM.WW.ScadaBridge.IntegrationTests.Cluster;
-
-///
-/// N1 regression (review 02 round 2, Critical): the resync authority must use the same
-/// oldest-Up predicate as the S&F delivery gate. Divergence scenario = the delivering node
-/// is OLDEST but not LEADER (leader = lowest address), the exact state a rolling restart of
-/// the lower-address node produces. Pre-fix the delivering node requests a resync from the
-/// stale peer and ReplaceAllAsync wipes its live buffer.
-///
-public class SfBufferResyncPredicateTests
-{
- [Fact]
- public async Task OldestButNotLeaderNode_KeepsItsBuffer_AndSeedsTheJoiner()
- {
- // Two explicit ports, deliberately assigned so the FIRST-started (oldest,
- // delivering) node has the HIGHER address → the second node is cluster leader.
- var p1 = TwoNodeClusterFixture.GetFreeTcpPort();
- var p2 = TwoNodeClusterFixture.GetFreeTcpPort();
- var (portHigh, portLow) = p1 > p2 ? (p1, p2) : (p2, p1);
-
- var fixture = await TwoNodeClusterFixture.StartAsync(
- role: "site-int", portA: portHigh, portB: portLow);
-
- // The S&F stores AND SiteStorageService take an ILocalDb (LocalDb has no in-memory
- // mode), so each node gets its own temp-file local databases. They are disposed
- // AFTER the cluster is shut down — the actors hold connections while the systems
- // are alive — and only then are the files (plus their WAL sidecars) deleted.
- var localDbs = new List();
-
- try
- {
- // Real S&F storage + replication actor per node, production default predicate
- // (no isActiveOverride) — the exact wiring under test.
- var (storageOldest, _, sfDbOldest, siteDbOldest) =
- await CreateReplicationActorAsync(fixture.NodeA, "oldest");
- localDbs.Add(sfDbOldest);
- localDbs.Add(siteDbOldest);
- var (storageJoiner, _, sfDbJoiner, siteDbJoiner) =
- await CreateReplicationActorAsync(fixture.NodeB, "joiner");
- localDbs.Add(sfDbJoiner);
- localDbs.Add(siteDbJoiner);
-
- // The delivering (oldest) node has a live buffered row the standby never saw.
- await storageOldest.EnqueueAsync(NewMessage("live-row"));
-
- // Trigger peer (re)tracking on both sides: each actor got InitialStateAsSnapshot
- // in PreStart, but the enqueue raced it — re-deliver via a fresh MemberUp is not
- // needed; OnPeerTracked already fired on join. The resync exchange is async:
- // wait until the JOINER holds the row (proves the snapshot flowed oldest→joiner,
- // the correct direction). Pre-fix this times out (the joiner, as leader, never
- // requests) AND the oldest node's row is deleted by the stale wipe.
- await AwaitAsync(async () => await storageJoiner.GetMessageByIdAsync("live-row") != null,
- TimeSpan.FromSeconds(20),
- "joiner never received the resync snapshot (resync ran in the wrong direction)");
-
- // And the delivering node's buffer is untouched — the N1 wipe assertion.
- Assert.NotNull(await storageOldest.GetMessageByIdAsync("live-row"));
- }
- finally
- {
- await fixture.DisposeAsync();
-
- foreach (var localDb in localDbs)
- {
- var path = localDb.Path;
- localDb.Dispose();
- TestLocalDb.DeleteFiles(path);
- }
- }
- }
-
- private static async Task<(
- StoreAndForwardStorage Storage, IActorRef Actor, TestLocalDb SfLocalDb, TestLocalDb SiteLocalDb)>
- CreateReplicationActorAsync(ActorSystem node, string tag)
- {
- var sfLocalDb = TestLocalDb.CreateTemp($"sf-resync-{tag}");
- var sfStorage = new StoreAndForwardStorage(sfLocalDb.Db,
- NullLogger.Instance);
- await sfStorage.InitializeAsync();
- var siteLocalDb = TestLocalDb.CreateTemp($"site-resync-{tag}");
- var siteStorage = new SiteStorageService(siteLocalDb.Db,
- NullLogger.Instance);
- var replicationService = new ReplicationService(
- new StoreAndForwardOptions(), NullLogger.Instance);
- // Name MUST be "site-replication" — SendToPeer targets /user/site-replication.
- var actor = node.ActorOf(Props.Create(() => new SiteReplicationActor(
- siteStorage, sfStorage, replicationService, "site-int",
- NullLogger.Instance, null, null, null, null)),
- "site-replication");
- return (sfStorage, actor, sfLocalDb, siteLocalDb);
- }
-
- private static StoreAndForwardMessage NewMessage(string id) => new()
- {
- Id = id,
- Category = StoreAndForwardCategory.Notification,
- Target = "central",
- PayloadJson = "{}",
- CreatedAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Pending,
- MaxRetries = 0,
- };
-
- private static async Task AwaitAsync(Func> condition, TimeSpan timeout, string why)
- {
- var deadline = DateTime.UtcNow + timeout;
- while (DateTime.UtcNow < deadline)
- {
- if (await condition()) return;
- await Task.Delay(250);
- }
- throw new TimeoutException(why);
- }
-}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairHarness.cs b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairHarness.cs
index dc37eefb..0bc73e80 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairHarness.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.IntegrationTests/LocalDbSitePairHarness.cs
@@ -129,40 +129,14 @@ public abstract class LocalDbSitePairHarness : IAsyncLifetime
.Build();
return new ServiceCollection()
- .AddZbLocalDb(config, db =>
- {
- SiteLocalDbSetup.OnReady(db, config);
- RegisterPhase2TablesUntilCutover(db);
- })
+ .AddZbLocalDb(config, db => SiteLocalDbSetup.OnReady(db, config))
.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.
+ /// The eight tables the Phase 2 cutover registers. Listed literally rather than derived
+ /// from production code, so a registration that drifts fails these tests instead of
+ /// agreeing with itself.
///
///
/// notification_lists and smtp_configurations are absent by design. They
diff --git a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerActorTests.cs b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerActorTests.cs
index f898add5..4f448aa0 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerActorTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerActorTests.cs
@@ -67,12 +67,14 @@ public class DeploymentManagerActorTests : TestKit, IDisposable
null, // no stream manager in tests
options,
NullLogger.Instance,
- null,
- null,
- null,
- serviceProvider,
- null,
- configFetcher)));
+ // Named from here on. These trailing parameters are all optional and several
+ // share a type, so a positional list silently binds the wrong argument when the
+ // signature changes — which is exactly what removing replicationActor did.
+ dclManager: null,
+ healthCollector: null,
+ serviceProvider: serviceProvider,
+ loggerFactory: null,
+ configFetcher: configFetcher)));
}
private static string MakeConfigJson(string instanceName)
@@ -301,8 +303,11 @@ public class DeploymentManagerActorTests : TestKit, IDisposable
var dm = ActorOf(Props.Create(() => new DeploymentManagerActor(
_storage, _compilationService, _sharedScriptLibrary, null,
new SiteRuntimeOptions(), NullLogger.Instance,
- null, null, null, null, null, null,
- TimeSpan.FromMilliseconds(200), loader)));
+ // dclManager, healthCollector, serviceProvider, loggerFactory, configFetcher.
+ // Props.Create builds an expression tree, which rejects named arguments that
+ // are out of position, so the optional tail has to be padded positionally.
+ null, null, null, null, null,
+ startupLoadRetryInterval: TimeSpan.FromMilliseconds(200), configLoader: loader)));
AwaitAssert(() =>
{
@@ -892,13 +897,14 @@ public class DeploymentManagerActorTests : TestKit, IDisposable
// Security cleanup. notification_lists and smtp_configurations can hold plaintext
// SMTP passwords written by a pre-2026-07-10 build, and the ACTIVE node's artifact
// apply is what clears them (DeploymentManagerActor.HandleDeployArtifacts). The
- // standby's copy of this call lives in SiteReplicationActor and dies with it at
- // Task 15, which makes this call site the ONLY remaining one.
+ // standby used to hold a second copy of this call in SiteReplicationActor; LocalDb
+ // Phase 2 deleted that actor, so this is now the ONLY call site that keeps the
+ // tables empty.
//
- // Nothing pins it today: ArtifactStorageTests covers the storage method, not the
- // actor's call to it, so Task 16's edits to this actor could drop the call and every
- // suite would stay green. That is precisely the kind of silent security regression
- // this test exists to prevent — verified red-first by commenting out the call.
+ // ArtifactStorageTests covers the storage method, not the actor's call to it, so
+ // without this test the call could be dropped and every suite would stay green.
+ // That is precisely the kind of silent security regression it exists to prevent —
+ // verified red-first by commenting out the call.
await SeedCentralOnlyRowsAsync();
Assert.Equal(1, await RowCountAsync("notification_lists"));
Assert.Equal(1, await RowCountAsync("smtp_configurations"));
diff --git a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerCertReconcileTests.cs b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerCertReconcileTests.cs
index 7ea67d09..11c6a0b4 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerCertReconcileTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerCertReconcileTests.cs
@@ -80,8 +80,7 @@ akka {
private IActorRef CreateDeploymentManager() =>
ActorOf(Props.Create(() => new DeploymentManagerActor(
_storage, _compilationService, _sharedScriptLibrary,
- null, new SiteRuntimeOptions(), NullLogger.Instance,
- null, null, null, null, null, null, null, null)));
+ null, new SiteRuntimeOptions(), NullLogger.Instance)));
[Fact]
public void SiteNodeJoined_PushesLocalTrustedCertsToJoinedNode()
diff --git a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerLoggerFactoryTests.cs b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerLoggerFactoryTests.cs
index 087b7f0d..fcf6da18 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerLoggerFactoryTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerLoggerFactoryTests.cs
@@ -103,11 +103,10 @@ public class DeploymentManagerLoggerFactoryTests : TestKit, IDisposable
null,
new SiteRuntimeOptions { StartupBatchSize = 100, StartupBatchDelayMs = 5 },
NullLogger.Instance,
- null,
- null,
- null,
- null,
- loggerFactory)));
+ // dclManager, healthCollector, serviceProvider — padded positionally because
+ // Props.Create is an expression tree and rejects out-of-position named args.
+ null, null, null,
+ loggerFactory: loggerFactory)));
// Allow async startup (load configs + staggered creation).
await Task.Delay(2000);
diff --git a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerRedeployTests.cs b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerRedeployTests.cs
index a853d061..663243a3 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerRedeployTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/DeploymentManagerRedeployTests.cs
@@ -61,10 +61,11 @@ public class DeploymentManagerRedeployTests : TestKit, IDisposable
null,
new SiteRuntimeOptions(),
NullLogger.Instance,
+ // dclManager — padded positionally because Props.Create is an expression tree
+ // and rejects out-of-position named args.
null,
- null,
- healthCollector,
- serviceProvider)));
+ healthCollector: healthCollector,
+ serviceProvider: serviceProvider)));
}
///
diff --git a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/SiteReplicationActorTests.cs b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/SiteReplicationActorTests.cs
deleted file mode 100644
index 0dae7d19..00000000
--- a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/Actors/SiteReplicationActorTests.cs
+++ /dev/null
@@ -1,568 +0,0 @@
-using System.Collections.Concurrent;
-using System.Diagnostics.Metrics;
-using Akka.Actor;
-using Akka.TestKit.Xunit2;
-using Microsoft.Extensions.Logging;
-using Microsoft.Extensions.Logging.Abstractions;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Deployment;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Messages;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Persistence;
-using ZB.MOM.WW.ScadaBridge.StoreAndForward;
-using ZB.MOM.WW.ScadaBridge.Commons.Observability;
-using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
-using ZB.MOM.WW.ScadaBridge.TestSupport;
-
-namespace ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests.Actors;
-
-///
-/// Tests for 's notify-and-fetch config replication:
-/// the active node now replicates an id-only (no inline
-/// config JSON — killing the intra-site 128 KB frame trap), and the standby fetches the
-/// config from central over HTTP and writes it with the older-write guard.
-///
-public class SiteReplicationActorTests : TestKit, IDisposable
-{
- // Cluster provider is required because SiteReplicationActor calls Cluster.Get in its ctor
- // and subscribes to cluster events in PreStart. We use the in-memory TestTransport (not
- // dot-netty) so no real socket is bound and no DNS lookup happens — the actor only needs
- // the cluster extension to load; these tests never form a real two-node cluster.
- private const string ClusterConfig = @"
-akka {
- actor { provider = cluster }
- remote {
- enabled-transports = [""akka.remote.test""]
- test {
- transport-class = ""Akka.Remote.Transport.TestTransport, Akka.Remote""
- applied-adapters = []
- registry-key = site-repl-test
- local-address = ""test://site-repl@localhost:1""
- maximum-payload-bytes = 128000b
- scheme-identifier = test
- }
- }
- cluster { roles = [""site-test""] }
- loglevel = WARNING
-}";
-
- private const string SiteRole = "site-test";
-
- private readonly SiteStorageService _storage;
- private readonly TestLocalDb _siteLocalDb;
- private readonly TestLocalDb _sfLocalDb;
- private readonly StoreAndForwardStorage _sfStorage;
- private readonly ReplicationService _replicationService;
- private readonly string _sfDbFile;
-
- public SiteReplicationActorTests() : base(ClusterConfig, "site-repl")
- {
- _sfDbFile = Path.Combine(Path.GetTempPath(), $"site-repl-sf-{Guid.NewGuid():N}.db");
-
- // SiteStorageService takes an ILocalDb now; LocalDb has no in-memory mode, so the
- // site store gets its own temp-file database alongside the S&F one.
- _siteLocalDb = TestLocalDb.CreateTemp("site-repl-test");
- _storage = new SiteStorageService(
- _siteLocalDb.Db, NullLogger.Instance);
- _storage.InitializeAsync().GetAwaiter().GetResult();
-
- _sfLocalDb = TestLocalDb.Create(_sfDbFile);
- _sfStorage = new StoreAndForwardStorage(
- _sfLocalDb.Db, NullLogger.Instance);
- _sfStorage.InitializeAsync().GetAwaiter().GetResult();
-
- _replicationService = new ReplicationService(
- new StoreAndForwardOptions(), NullLogger.Instance);
- }
-
- void IDisposable.Dispose()
- {
- Shutdown();
- // The master connection anchors the WAL — dispose before deleting.
- var siteDbPath = _siteLocalDb.Path;
- _siteLocalDb.Dispose();
- _sfLocalDb.Dispose();
- TestLocalDb.DeleteFiles(siteDbPath);
- TestLocalDb.DeleteFiles(_sfDbFile);
- }
-
- private IActorRef CreateReplicationActor(IDeploymentConfigFetcher fetcher) =>
- ActorOf(Props.Create(() => new SiteReplicationActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, fetcher)));
-
- private IActorRef CreateReplicationActor(
- IDeploymentConfigFetcher fetcher, SiteRuntimeOptions options, TimeSpan retryDelay) =>
- ActorOf(Props.Create(() => new SiteReplicationActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, fetcher, null, options, retryDelay)));
-
- [Fact]
- public async Task ReplicatedFetch_RetriesUpToConfigFetchRetryCount()
- {
- // The first two fetches fail transiently; the third succeeds. With
- // ConfigFetchRetryCount = 3 the standby must retry to the third attempt and
- // then guarded-write the fetched config (a short retry delay keeps the test fast).
- var attempts = 0;
- var fetcher = new FakeConfigFetcher(_ =>
- Interlocked.Increment(ref attempts) < 3
- ? Task.FromException(new InvalidOperationException("central hiccup"))
- : Task.FromResult("{\"instanceUniqueName\":\"RetryPump\"}"));
- var actor = CreateReplicationActor(
- fetcher, new SiteRuntimeOptions { ConfigFetchRetryCount = 3 },
- TimeSpan.FromMilliseconds(50));
-
- actor.Tell(new ApplyConfigDeploy(
- "RetryPump", "dep-r1", "sha256:r1", true,
- "http://central:9000", "tok-r1"));
-
- await AwaitAssertAsync(async () =>
- {
- Assert.Equal(3, Volatile.Read(ref attempts));
- var configs = await _storage.GetAllDeployedConfigsAsync();
- Assert.Single(configs, c => c.InstanceUniqueName == "RetryPump");
- }, TimeSpan.FromSeconds(10));
- }
-
- [Fact]
- public async Task ApplyConfigDeploy_StandbyFetchesConfigAndGuardedWrites()
- {
- // The standby receives an id-only ApplyConfigDeploy; it fetches the config from
- // central using the message's coords, then guarded-writes the fetched config.
- const string configJson = "{\"instanceUniqueName\":\"Pump1\"}";
- var fetcher = new FakeConfigFetcher(_ => Task.FromResult(configJson));
- var actor = CreateReplicationActor(fetcher);
-
- actor.Tell(new ApplyConfigDeploy(
- "Pump1", "dep-100", "sha256:abc", true,
- "http://central:9000", "tok-xyz"));
-
- // The continuation runs off-thread; await the guarded write landing.
- await AwaitAssertAsync(async () =>
- {
- var configs = await _storage.GetAllDeployedConfigsAsync();
- var row = Assert.Single(configs, c => c.InstanceUniqueName == "Pump1");
- Assert.Equal(configJson, row.ConfigJson);
- Assert.Equal("dep-100", row.DeploymentId);
- Assert.Equal("sha256:abc", row.RevisionHash);
- Assert.True(row.IsEnabled);
- }, TimeSpan.FromSeconds(5));
-
- // The fetcher was called with the message's coords.
- var call = Assert.Single(fetcher.Calls);
- Assert.Equal("http://central:9000", call.BaseUrl);
- Assert.Equal("dep-100", call.DeploymentId);
- Assert.Equal("tok-xyz", call.Token);
- }
-
- [Fact]
- public async Task ApplyConfigDeploy_Superseded404_SkipsWriteAndActorSurvives()
- {
- // A 404 (superseded/expired) surfaces as DeploymentConfigFetchException{IsSuperseded}.
- // The standby must skip the write, observe the exception (no crash), and stay alive.
- var fetcher = new FakeConfigFetcher(_ =>
- Task.FromException(
- new DeploymentConfigFetchException("expired", isSuperseded: true)));
- var actor = CreateReplicationActor(fetcher);
-
- actor.Tell(new ApplyConfigDeploy(
- "GonePump", "dep-stale", "sha256:gone", true,
- "http://central:9000", "tok-stale"));
-
- // The fetch was attempted...
- await AwaitAssertAsync(() =>
- {
- Assert.Single(fetcher.Calls);
- return Task.CompletedTask;
- }, TimeSpan.FromSeconds(5));
-
- // ...the actor did not crash (no Terminated to its watcher within the window)...
- Watch(actor);
- ExpectNoMsg(TimeSpan.FromMilliseconds(500));
-
- // ...and nothing was written for the superseded instance.
- var configs = await _storage.GetAllDeployedConfigsAsync();
- Assert.DoesNotContain(configs, c => c.InstanceUniqueName == "GonePump");
- }
-
- [Fact]
- public async Task ApplyConfigDeploy_EmptyFetchCoords_SkipsFetchAndWrite()
- {
- // The direct DeployInstanceCommand cross-cluster wire path was retired in Task 14.
- // This tests the defensive guard: if empty coords arrive, the actor must skip quietly
- // — no FetchAsync("") call, no write — rather than erroring.
- var fetcher = new FakeConfigFetcher(_ => Task.FromResult("never"));
- var actor = CreateReplicationActor(fetcher);
-
- actor.Tell(new ApplyConfigDeploy(
- "NoCoordsPump", "dep-direct", "sha256:nc", true,
- CentralFetchBaseUrl: "", FetchToken: ""));
-
- // Give any (erroneous) async continuation time to run, then prove neither happened.
- Watch(actor);
- ExpectNoMsg(TimeSpan.FromMilliseconds(500));
- Assert.Empty(fetcher.Calls);
- var configs = await _storage.GetAllDeployedConfigsAsync();
- Assert.DoesNotContain(configs, c => c.InstanceUniqueName == "NoCoordsPump");
- }
-
- [Fact]
- public void ReplicateConfigDeploy_MapsToIdOnlyApplyConfigDeploy_ForPeer()
- {
- // The outbound mapping must forward an id-only ApplyConfigDeploy carrying the fetch
- // coords (and NO inline config) to the peer.
- var probe = CreateTestProbe();
- var fetcher = new FakeConfigFetcher(_ => Task.FromResult("unused"));
- var actor = ActorOf(Props.Create(() => new ProbeForwardingReplicationActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, fetcher, probe.Ref)));
-
- actor.Tell(new ReplicateConfigDeploy(
- "Pump2", "dep-200", "sha256:def", false,
- "http://central:9000", "tok-abc"));
-
- var applied = probe.ExpectMsg(TimeSpan.FromSeconds(3));
- Assert.Equal("Pump2", applied.InstanceName);
- Assert.Equal("dep-200", applied.DeploymentId);
- Assert.Equal("sha256:def", applied.RevisionHash);
- Assert.False(applied.IsEnabled);
- Assert.Equal("http://central:9000", applied.CentralFetchBaseUrl);
- Assert.Equal("tok-abc", applied.FetchToken);
- }
-
- // ── Task 21: peer-join S&F buffer resync (anti-entropy) ──
-
- [Fact]
- public void StandbyTrackingPeer_SendsResyncRequest()
- {
- var probe = CreateTestProbe();
- var actor = ActorOf(Props.Create(() => new ResyncTestActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, probe.Ref, () => false)));
-
- actor.Tell(new TriggerPeerTracked()); // stands in for TryTrackPeer's MemberUp path
-
- probe.ExpectMsg(TimeSpan.FromSeconds(3));
- }
-
- [Fact]
- public void ActiveTrackingPeer_DoesNotRequestResync()
- {
- var probe = CreateTestProbe();
- var actor = ActorOf(Props.Create(() => new ResyncTestActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, probe.Ref, () => true)));
-
- actor.Tell(new TriggerPeerTracked());
-
- probe.ExpectNoMsg(TimeSpan.FromMilliseconds(300)); // active node never requests a resync
- }
-
- [Fact]
- public async Task ActiveNode_AnswersResyncRequest_WithChunkedSnapshot()
- {
- // Post-R2-T5 the active node answers with byte-budgeted SfBufferSnapshotChunk(s)
- // (a single small row rides one chunk) rather than the monolithic SfBufferSnapshot.
- await _sfStorage.EnqueueAsync(NewSfMessage("m1"));
- var probe = CreateTestProbe();
- var actor = ActorOf(Props.Create(() => new ResyncTestActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, probe.Ref, () => true)));
-
- actor.Tell(new RequestSfBufferResync(), TestActor);
-
- var chunk = ExpectMsg(TimeSpan.FromSeconds(3));
- Assert.Equal(1, chunk.TotalChunks);
- Assert.Equal(1, chunk.Sequence);
- Assert.Single(chunk.Messages);
- Assert.False(chunk.Truncated);
- }
-
- [Fact]
- public async Task StandbyNode_AppliesSnapshot_ReplacingItsBuffer()
- {
- await _sfStorage.EnqueueAsync(NewSfMessage("stale"));
- var probe = CreateTestProbe();
- var actor = ActorOf(Props.Create(() => new ResyncTestActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, probe.Ref, () => false)));
-
- actor.Tell(new SfBufferSnapshot(new List { NewSfMessage("fresh") }, false));
-
- await AwaitAssertAsync(async () =>
- {
- Assert.Null(await _sfStorage.GetMessageByIdAsync("stale"));
- Assert.NotNull(await _sfStorage.GetMessageByIdAsync("fresh"));
- }, TimeSpan.FromSeconds(5));
- }
-
- // ── R2 T5: chunked resync answer ──
-
- [Fact]
- public void ChunkForRemoting_SplitsByByteBudget_PreservingOrderAndSequence()
- {
- var rows = Enumerable.Range(0, 10)
- .Select(i => NewMessage($"m{i}", payloadJson: new string('x', 20_000)))
- .ToList();
-
- var chunks = SiteReplicationActor.ChunkForRemoting(rows, maxChunkBytes: 64_000, maxChunkRows: 200);
-
- Assert.True(chunks.Count > 1); // 10 × 20 KB cannot ride one 64 KB chunk
- Assert.Equal(rows.Select(r => r.Id), chunks.SelectMany(c => c).Select(r => r.Id)); // order preserved
- Assert.All(chunks, c => Assert.True(
- c.Sum(r => r.PayloadJson.Length) <= 64_000 || c.Count == 1)); // budget honored (oversized row isolated)
- }
-
- [Fact]
- public void ChunkForRemoting_RowCapHonored_AndSingleOversizedRowIsolated()
- {
- var many = Enumerable.Range(0, 500).Select(i => NewMessage($"s{i}", payloadJson: "{}")).ToList();
- Assert.All(SiteReplicationActor.ChunkForRemoting(many, 64_000, 200), c => Assert.True(c.Count <= 200));
-
- var oversized = new List
- { NewMessage("big", payloadJson: new string('y', 100_000)), NewMessage("small", payloadJson: "{}") };
- var chunks = SiteReplicationActor.ChunkForRemoting(oversized, 64_000, 200);
- Assert.Equal(2, chunks.Count); // the oversized row rides alone
- }
-
- [Fact]
- public async Task ActiveNode_AnswersResyncRequest_WithSequencedChunks_SharingOneResyncId()
- {
- for (var i = 0; i < 3; i++)
- await _sfStorage.EnqueueAsync(NewMessage($"c{i}", payloadJson: new string('z', 30_000)));
- var actor = CreateResyncActor(isActive: () => true);
-
- actor.Tell(new RequestSfBufferResync(), TestActor);
-
- var first = ExpectMsg(TimeSpan.FromSeconds(5));
- var rest = Enumerable.Range(1, first.TotalChunks - 1)
- .Select(_ => ExpectMsg(TimeSpan.FromSeconds(5)))
- .Prepend(first)
- .ToList();
-
- Assert.True(first.TotalChunks > 1);
- Assert.All(rest, c => Assert.Equal(first.ResyncId, c.ResyncId));
- Assert.Equal(Enumerable.Range(1, first.TotalChunks), rest.Select(c => c.Sequence));
- Assert.Equal(3, rest.Sum(c => c.Messages.Count));
- }
-
- // ── R2 T6: standby chunk assembly + atomic apply + ack ──
-
- [Fact]
- public async Task StandbyNode_AssemblesChunks_AppliesOnce_AndAcks()
- {
- await _sfStorage.EnqueueAsync(NewMessage("stale"));
- var actor = CreateResyncActor(isActive: () => false);
- var resyncId = "r1";
-
- actor.Tell(new SfBufferSnapshotChunk(resyncId, 1, 2,
- new List { NewMessage("f1") }, false), TestActor);
- actor.Tell(new SfBufferSnapshotChunk(resyncId, 2, 2,
- new List { NewMessage("f2") }, false), TestActor);
-
- var ack = ExpectMsg(TimeSpan.FromSeconds(5));
- Assert.Equal(resyncId, ack.ResyncId);
- Assert.Equal(2, ack.RowCount);
- await AwaitAssertAsync(async () =>
- {
- Assert.Null(await _sfStorage.GetMessageByIdAsync("stale")); // replaced wholesale
- Assert.NotNull(await _sfStorage.GetMessageByIdAsync("f1"));
- Assert.NotNull(await _sfStorage.GetMessageByIdAsync("f2"));
- });
- }
-
- [Fact]
- public async Task StandbyNode_NewResyncId_DiscardsStalePartialAssembly()
- {
- var actor = CreateResyncActor(isActive: () => false);
- actor.Tell(new SfBufferSnapshotChunk("old", 1, 2,
- new List { NewMessage("orphan") }, false), TestActor);
- actor.Tell(new SfBufferSnapshotChunk("new", 1, 1,
- new List { NewMessage("fresh") }, false), TestActor);
-
- ExpectMsg(TimeSpan.FromSeconds(5)); // "new" completed
- await AwaitAssertAsync(async () =>
- {
- Assert.NotNull(await _sfStorage.GetMessageByIdAsync("fresh"));
- Assert.Null(await _sfStorage.GetMessageByIdAsync("orphan")); // stale partial never applied
- });
- }
-
- [Fact]
- public void ActiveNode_IgnoresChunks_NeverAcks()
- {
- var actor = CreateResyncActor(isActive: () => true);
- actor.Tell(new SfBufferSnapshotChunk("r", 1, 1,
- new List { NewMessage("x") }, false), TestActor);
- ExpectNoMsg(TimeSpan.FromMilliseconds(300));
- }
-
- // ── R2 T7: active-side resync ack confirmation + telemetry ──
- //
- // NOTE (deviation from plan): the actor logs via Microsoft ILogger (NullLogger in
- // tests), NOT Akka's EventStream, so the plan's EventFilter.Warning assertions can
- // never observe these warnings. We observe the two OTel counters via a MeterListener
- // instead — the equivalent, and stronger, observable signal.
-
- [Fact]
- public async Task ActiveNode_ReceivingAck_CountsResyncCompleted()
- {
- long completed = 0;
- using var listener = ListenCounter("scadabridge.store_and_forward.resync.completed",
- m => Interlocked.Add(ref completed, m));
-
- await _sfStorage.EnqueueAsync(NewMessage("m1"));
- var actor = CreateResyncActor(isActive: () => true, ackTimeout: TimeSpan.FromSeconds(30));
- actor.Tell(new RequestSfBufferResync(), TestActor);
- var chunk = ExpectMsg(TimeSpan.FromSeconds(5));
-
- actor.Tell(new SfBufferResyncAck(chunk.ResyncId, 1), TestActor);
-
- await AwaitAssertAsync(() =>
- {
- Assert.True(Interlocked.Read(ref completed) >= 1); // ack recorded the completion
- return Task.CompletedTask;
- }, TimeSpan.FromSeconds(5));
- }
-
- [Fact]
- public async Task ActiveNode_MissingAck_WarnsAfterAckTimeout()
- {
- long ackMissing = 0;
- using var listener = ListenCounter("scadabridge.store_and_forward.resync.ack_missing",
- m => Interlocked.Add(ref ackMissing, m));
-
- await _sfStorage.EnqueueAsync(NewMessage("m1"));
- var actor = CreateResyncActor(isActive: () => true, ackTimeout: TimeSpan.FromMilliseconds(200));
- actor.Tell(new RequestSfBufferResync(), TestActor);
- ExpectMsg(TimeSpan.FromSeconds(5));
- // No ack is sent → the ack window expires and the resync is counted unacknowledged.
-
- await AwaitAssertAsync(() =>
- {
- Assert.True(Interlocked.Read(ref ackMissing) >= 1);
- return Task.CompletedTask;
- }, TimeSpan.FromSeconds(5));
- }
-
- /// Attaches a to a single ScadaBridge counter by name,
- /// forwarding each recorded increment to .
- private static MeterListener ListenCounter(string instrumentName, Action onMeasurement)
- {
- var listener = new MeterListener();
- listener.InstrumentPublished = (inst, l) =>
- {
- if (inst.Meter.Name == ScadaBridgeTelemetry.MeterName && inst.Name == instrumentName)
- l.EnableMeasurementEvents(inst);
- };
- listener.SetMeasurementEventCallback((_, m, _, _) => onMeasurement(m));
- listener.Start();
- return listener;
- }
-
- private static StoreAndForwardMessage NewSfMessage(string id) => new()
- {
- Id = id,
- Category = StoreAndForwardCategory.ExternalSystem,
- Target = "t",
- PayloadJson = "{}",
- RetryCount = 0,
- MaxRetries = 50,
- RetryIntervalMs = 30000,
- CreatedAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Pending,
- };
-
- ///
- /// Builds a resync-test message with a settable payload (additive to
- /// — the chunker sizes on PayloadJson length).
- ///
- private static StoreAndForwardMessage NewMessage(string id, string payloadJson = "{}") => new()
- {
- Id = id,
- Category = StoreAndForwardCategory.ExternalSystem,
- Target = "t",
- PayloadJson = payloadJson,
- RetryCount = 0,
- MaxRetries = 50,
- RetryIntervalMs = 30000,
- CreatedAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Pending,
- };
-
- /// Constructs a with the given active-node check
- /// (the resync chunk/ack tests Tell to and expect from ).
- /// is the active-side ack window seam (T7).
- private IActorRef CreateResyncActor(Func isActive, TimeSpan? ackTimeout = null) =>
- ActorOf(Props.Create(() => new ResyncTestActor(
- _storage, _sfStorage, _replicationService, SiteRole,
- NullLogger.Instance, CreateTestProbe().Ref, isActive, ackTimeout)));
-
- /// Test message: drives directly,
- /// standing in for the MemberUp→TryTrackPeer path (a single-node TestKit cannot form a real peer).
- private sealed record TriggerPeerTracked;
-
- ///
- /// Test subclass for the resync tests: captures peer sends to a probe, injects the
- /// active-node check, and exposes via a test message.
- ///
- private sealed class ResyncTestActor : SiteReplicationActor
- {
- private readonly IActorRef _peerProbe;
-
- public ResyncTestActor(
- SiteStorageService storage, StoreAndForwardStorage sfStorage,
- ReplicationService replicationService, string siteRole,
- ILogger logger, IActorRef peerProbe, Func isActive,
- TimeSpan? ackTimeout = null)
- : base(storage, sfStorage, replicationService, siteRole, logger,
- configFetcher: null, isActiveOverride: isActive, resyncAckTimeout: ackTimeout)
- {
- _peerProbe = peerProbe;
- Receive(_ => OnPeerTracked());
- }
-
- protected override void SendToPeer(object message) => _peerProbe.Tell(message, Self);
- }
-
- ///
- /// Test subclass exposing the peer send: is
- /// overridden to forward to a probe so the outbound mapping can be asserted without a real
- /// two-node cluster (a single-node TestKit has no peer address, so the real send is dropped).
- ///
- private sealed class ProbeForwardingReplicationActor : SiteReplicationActor
- {
- private readonly IActorRef _peerProbe;
-
- public ProbeForwardingReplicationActor(
- SiteStorageService storage, StoreAndForwardStorage sfStorage,
- ReplicationService replicationService, string siteRole,
- ILogger logger, IDeploymentConfigFetcher configFetcher,
- IActorRef peerProbe)
- : base(storage, sfStorage, replicationService, siteRole, logger, configFetcher)
- => _peerProbe = peerProbe;
-
- protected override void SendToPeer(object message) => _peerProbe.Tell(message, Self);
- }
-
- ///
- /// In-test fake : runs a per-deploymentId behavior
- /// (return config JSON or throw, as a Task — mirroring the real async HTTP fetcher) and
- /// records every call's coords thread-safely (the continuation runs on a pool thread).
- ///
- private sealed class FakeConfigFetcher : IDeploymentConfigFetcher
- {
- private readonly Func> _behavior;
- public ConcurrentQueue<(string BaseUrl, string DeploymentId, string Token)> Calls { get; } = new();
-
- public FakeConfigFetcher(Func> behavior) => _behavior = behavior;
-
- public async Task FetchAsync(
- string centralFetchBaseUrl, string deploymentId, string token, CancellationToken ct)
- {
- Calls.Enqueue((centralFetchBaseUrl, deploymentId, token));
- await Task.Yield();
- return await _behavior(deploymentId);
- }
- }
-}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/ResyncWireSerializationPinTests.cs b/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/ResyncWireSerializationPinTests.cs
deleted file mode 100644
index 43560a14..00000000
--- a/tests/ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests/ResyncWireSerializationPinTests.cs
+++ /dev/null
@@ -1,95 +0,0 @@
-using Akka.TestKit.Xunit2;
-using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
-using ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors;
-using ZB.MOM.WW.ScadaBridge.StoreAndForward;
-
-namespace ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests;
-
-///
-/// Characterization pin for the chunked anti-entropy resync contract (review 02 round 2,
-/// N2). (active→standby) and
-/// (standby→active) ride intra-site Akka remoting on the default reflective-JSON wire
-/// format. A rename/move, dropped setter, or non-default-constructible message would
-/// silently break resync across a rolling upgrade and only surface as a divergent buffer
-/// after a failover. These pin round-trip fidelity and type identity. (These messages are
-/// NOT ClusterClient traffic, so they are intentionally absent from ClusterClientContractLockTests.)
-///
-public class ResyncWireSerializationPinTests : TestKit
-{
- private T RoundTrip(T message)
- {
- var serialization = Sys.Serialization;
- var serializer = serialization.FindSerializerFor(message);
- var bytes = serializer.ToBinary(message);
- return (T)serialization.Deserialize(bytes, serializer.Identifier, message!.GetType());
- }
-
- private static StoreAndForwardMessage FullMessage() => new()
- {
- Id = Guid.NewGuid().ToString("N"),
- Category = StoreAndForwardCategory.Notification,
- Target = "Operators",
- PayloadJson = "{\"notificationId\":\"abc\"}",
- RetryCount = 4,
- MaxRetries = 0,
- RetryIntervalMs = 30000,
- CreatedAt = DateTimeOffset.UtcNow,
- LastAttemptAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Parked,
- LastError = "central rejected",
- OriginInstanceName = "Plant.Pump3",
- ExecutionId = Guid.NewGuid(),
- SourceScript = "ScriptActor:MonitorSpeed",
- ParentExecutionId = Guid.NewGuid(),
- };
-
- [Fact]
- public void SfBufferSnapshotChunk_WithFullMessage_RoundTripsOnTheWire()
- {
- var message = FullMessage();
- var original = new SfBufferSnapshotChunk(
- "resync-1", 2, 5, new List { message }, Truncated: true);
-
- var back = RoundTrip(original);
-
- Assert.Equal(original.ResyncId, back.ResyncId);
- Assert.Equal(original.Sequence, back.Sequence);
- Assert.Equal(original.TotalChunks, back.TotalChunks);
- Assert.Equal(original.Truncated, back.Truncated);
- var m = Assert.Single(back.Messages);
- Assert.Equal(message.Id, m.Id);
- Assert.Equal(message.Category, m.Category);
- Assert.Equal(message.Target, m.Target);
- Assert.Equal(message.PayloadJson, m.PayloadJson);
- Assert.Equal(message.RetryCount, m.RetryCount);
- Assert.Equal(message.MaxRetries, m.MaxRetries);
- Assert.Equal(message.RetryIntervalMs, m.RetryIntervalMs);
- Assert.Equal(message.CreatedAt, m.CreatedAt);
- Assert.Equal(message.LastAttemptAt, m.LastAttemptAt);
- Assert.Equal(message.Status, m.Status);
- Assert.Equal(message.LastError, m.LastError);
- Assert.Equal(message.OriginInstanceName, m.OriginInstanceName);
- Assert.Equal(message.ExecutionId, m.ExecutionId);
- Assert.Equal(message.SourceScript, m.SourceScript);
- Assert.Equal(message.ParentExecutionId, m.ParentExecutionId);
- }
-
- [Fact]
- public void SfBufferResyncAck_RoundTripsOnTheWire()
- {
- var original = new SfBufferResyncAck("resync-1", 42);
-
- var back = RoundTrip(original);
-
- Assert.Equal(original.ResyncId, back.ResyncId);
- Assert.Equal(original.RowCount, back.RowCount);
- }
-
- // Type-identity pins: the reflective-JSON wire embeds CLR type manifests, so a
- // rename/move of either type silently breaks resync across a rolling upgrade.
- [Theory]
- [InlineData(typeof(SfBufferSnapshotChunk), "ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors.SfBufferSnapshotChunk")]
- [InlineData(typeof(SfBufferResyncAck), "ZB.MOM.WW.ScadaBridge.SiteRuntime.Actors.SfBufferResyncAck")]
- public void ResyncContract_TypeIdentity_IsPinned(Type type, string expectedFullName) =>
- Assert.Equal(expectedFullName, type.FullName);
-}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/CachedCallAttemptEmissionTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/CachedCallAttemptEmissionTests.cs
index 4bf217e7..f01746fa 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/CachedCallAttemptEmissionTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/CachedCallAttemptEmissionTests.cs
@@ -45,7 +45,6 @@ public class CachedCallAttemptEmissionTests : IAsyncLifetime, IDisposable
_storage,
_options,
NullLogger.Instance,
- replication: null,
cachedCallObserver: _observer,
siteId: "site-77");
}
@@ -490,7 +489,6 @@ public class CachedCallAttemptEmissionTests : IAsyncLifetime, IDisposable
RetryTimerInterval = TimeSpan.FromHours(1), // timer never fires in-test
},
NullLogger.Instance,
- replication: null,
cachedCallObserver: observer,
siteId: "site-77");
diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/ReplicationServiceTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/ReplicationServiceTests.cs
deleted file mode 100644
index 19fa4b2b..00000000
--- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/ReplicationServiceTests.cs
+++ /dev/null
@@ -1,236 +0,0 @@
-using Microsoft.Extensions.Logging;
-using Microsoft.Extensions.Logging.Abstractions;
-using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
-using ZB.MOM.WW.ScadaBridge.TestSupport;
-
-namespace ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests;
-
-///
-/// WP-11: Tests for async replication to standby.
-///
-public class ReplicationServiceTests : IAsyncLifetime, IDisposable
-{
- private readonly TestLocalDb _localDb;
- private readonly StoreAndForwardStorage _storage;
- private readonly ReplicationService _replicationService;
-
- public ReplicationServiceTests()
- {
- _localDb = TestLocalDb.CreateTemp("RepTests");
-
- _storage = new StoreAndForwardStorage(_localDb.Db, NullLogger.Instance);
-
- var options = new StoreAndForwardOptions { ReplicationEnabled = true };
- _replicationService = new ReplicationService(
- options, NullLogger.Instance);
- }
-
- public async Task InitializeAsync() => await _storage.InitializeAsync();
-
- public Task DisposeAsync() => Task.CompletedTask;
-
- public void Dispose()
- {
- var path = _localDb.Path;
- _localDb.Dispose();
- TestLocalDb.DeleteFiles(path);
- }
-
- [Fact]
- public void ReplicateEnqueue_NoHandler_DoesNotThrow()
- {
- var msg = CreateMessage("rep1");
- _replicationService.ReplicateEnqueue(msg);
- }
-
- [Fact]
- public async Task ReplicateEnqueue_WithHandler_ForwardsOperation()
- {
- ReplicationOperation? captured = null;
- _replicationService.SetReplicationHandler(op =>
- {
- captured = op;
- return Task.CompletedTask;
- });
-
- var msg = CreateMessage("rep2");
- _replicationService.ReplicateEnqueue(msg);
-
- await Task.Delay(200);
-
- Assert.NotNull(captured);
- Assert.Equal(ReplicationOperationType.Add, captured!.OperationType);
- Assert.Equal("rep2", captured.MessageId);
- }
-
- [Fact]
- public async Task ReplicateRemove_WithHandler_ForwardsRemoveOperation()
- {
- ReplicationOperation? captured = null;
- _replicationService.SetReplicationHandler(op =>
- {
- captured = op;
- return Task.CompletedTask;
- });
-
- _replicationService.ReplicateRemove("rep3");
-
- await Task.Delay(200);
-
- Assert.NotNull(captured);
- Assert.Equal(ReplicationOperationType.Remove, captured!.OperationType);
- Assert.Equal("rep3", captured.MessageId);
- }
-
- [Fact]
- public async Task ReplicatePark_WithHandler_ForwardsParkOperation()
- {
- ReplicationOperation? captured = null;
- _replicationService.SetReplicationHandler(op =>
- {
- captured = op;
- return Task.CompletedTask;
- });
-
- var msg = CreateMessage("rep4");
- _replicationService.ReplicatePark(msg);
-
- await Task.Delay(200);
-
- Assert.NotNull(captured);
- Assert.Equal(ReplicationOperationType.Park, captured!.OperationType);
- }
-
- [Fact]
- public async Task ApplyReplicatedOperationAsync_Add_EnqueuesMessage()
- {
- var msg = CreateMessage("apply1");
- var operation = new ReplicationOperation(ReplicationOperationType.Add, "apply1", msg);
-
- await _replicationService.ApplyReplicatedOperationAsync(operation, _storage);
-
- var retrieved = await _storage.GetMessageByIdAsync("apply1");
- Assert.NotNull(retrieved);
- }
-
- [Fact]
- public async Task ApplyReplicatedOperationAsync_Remove_DeletesMessage()
- {
- var msg = CreateMessage("apply2");
- await _storage.EnqueueAsync(msg);
-
- var operation = new ReplicationOperation(ReplicationOperationType.Remove, "apply2", null);
- await _replicationService.ApplyReplicatedOperationAsync(operation, _storage);
-
- var retrieved = await _storage.GetMessageByIdAsync("apply2");
- Assert.Null(retrieved);
- }
-
- [Fact]
- public async Task ApplyReplicatedOperationAsync_Park_UpdatesStatus()
- {
- var msg = CreateMessage("apply3");
- await _storage.EnqueueAsync(msg);
-
- var operation = new ReplicationOperation(ReplicationOperationType.Park, "apply3", msg);
- await _replicationService.ApplyReplicatedOperationAsync(operation, _storage);
-
- var retrieved = await _storage.GetMessageByIdAsync("apply3");
- Assert.NotNull(retrieved);
- Assert.Equal(StoreAndForwardMessageStatus.Parked, retrieved!.Status);
- }
-
- [Fact]
- public void ReplicateEnqueue_WhenReplicationDisabled_DoesNothing()
- {
- var options = new StoreAndForwardOptions { ReplicationEnabled = false };
- var service = new ReplicationService(options, NullLogger.Instance);
-
- bool handlerCalled = false;
- service.SetReplicationHandler(_ => { handlerCalled = true; return Task.CompletedTask; });
-
- service.ReplicateEnqueue(CreateMessage("disabled1"));
-
- Assert.False(handlerCalled);
- }
-
- [Fact]
- public async Task ReplicateEnqueue_HandlerThrows_DoesNotPropagateException()
- {
- _replicationService.SetReplicationHandler(_ =>
- throw new InvalidOperationException("standby down"));
-
- _replicationService.ReplicateEnqueue(CreateMessage("err1"));
-
- await Task.Delay(200);
- // No exception -- fire-and-forget, best-effort
- }
-
- // ── Task 10 (arch review 02): ordered, observable replication dispatch ──
-
- [Fact]
- public void ReplicationOperations_AreDispatchedInIssueOrder()
- {
- var seen = new List<(ReplicationOperationType, string)>();
- _replicationService.SetReplicationHandler(op =>
- {
- seen.Add((op.OperationType, op.MessageId));
- return Task.CompletedTask;
- });
-
- for (var i = 0; i < 200; i++)
- {
- _replicationService.ReplicateEnqueue(CreateMessage($"m{i}"));
- _replicationService.ReplicateRemove($"m{i}");
- }
-
- // Inline dispatch: by the time the calls return, every op was handed to the
- // handler, Add strictly before Remove per id. Pre-fix (Task.Run) this was
- // racy in both count and order.
- Assert.Equal(400, seen.Count);
- for (var i = 0; i < 200; i++)
- {
- Assert.Equal((ReplicationOperationType.Add, $"m{i}"), seen[2 * i]);
- Assert.Equal((ReplicationOperationType.Remove, $"m{i}"), seen[2 * i + 1]);
- }
- }
-
- [Fact]
- public void ReplicationHandlerThrow_IsSwallowed_AndLoggedAtWarning()
- {
- var logger = new CapturingLogger();
- var service = new ReplicationService(new StoreAndForwardOptions(), logger);
- service.SetReplicationHandler(_ => throw new InvalidOperationException("peer gone"));
-
- service.ReplicateRemove("m1"); // must not throw
-
- Assert.Contains(logger.Entries, e => e.Level == LogLevel.Warning);
- }
-
- private static StoreAndForwardMessage CreateMessage(string id)
- {
- return new StoreAndForwardMessage
- {
- Id = id,
- Category = StoreAndForwardCategory.ExternalSystem,
- Target = "target",
- PayloadJson = "{}",
- RetryCount = 0,
- MaxRetries = 50,
- RetryIntervalMs = 30000,
- CreatedAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Pending
- };
- }
-
- /// Minimal in-memory logger that records level + rendered message.
- private sealed class CapturingLogger : ILogger
- {
- public List<(LogLevel Level, string Message)> Entries { get; } = new();
- public IDisposable? BeginScope(TState state) where TState : notnull => null;
- public bool IsEnabled(LogLevel logLevel) => true;
- public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception,
- Func formatter)
- => Entries.Add((logLevel, formatter(state, exception)));
- }
-}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/ReplicationWireSerializationPinTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/ReplicationWireSerializationPinTests.cs
deleted file mode 100644
index f1d194f0..00000000
--- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/ReplicationWireSerializationPinTests.cs
+++ /dev/null
@@ -1,82 +0,0 @@
-using Akka.TestKit.Xunit2;
-using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
-
-namespace ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests;
-
-///
-/// Characterization pin for the intra-cluster S&F replication contract. A
-/// (carrying a full )
-/// is Told from the active node to the standby node's replication actor and rides
-/// the Akka default reflective-JSON wire format. Nothing tested that it round-trips
-/// or that its type identity is stable — a rename/move, a dropped setter, or a
-/// non-default-constructible message would silently break standby buffer sync and
-/// only surface as a divergent buffer after a failover.
-///
-/// The serializer swap (proto / explicit bindings) is deferred; until then this pin
-/// is the guardrail.
-///
-public class ReplicationWireSerializationPinTests : TestKit
-{
- private T RoundTrip(T message)
- {
- var serialization = Sys.Serialization;
- var serializer = serialization.FindSerializerFor(message);
- var bytes = serializer.ToBinary(message);
- return (T)serialization.Deserialize(bytes, serializer.Identifier, message!.GetType());
- }
-
- [Fact]
- public void ReplicationOperation_WithFullMessage_RoundTripsOnTheWire()
- {
- var message = new StoreAndForwardMessage
- {
- Id = Guid.NewGuid().ToString("N"),
- Category = StoreAndForwardCategory.Notification,
- Target = "Operators",
- PayloadJson = "{\"notificationId\":\"abc\"}",
- RetryCount = 4,
- MaxRetries = 0,
- RetryIntervalMs = 30000,
- CreatedAt = DateTimeOffset.UtcNow,
- LastAttemptAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Parked,
- LastError = "central rejected",
- OriginInstanceName = "Plant.Pump3",
- ExecutionId = Guid.NewGuid(),
- SourceScript = "ScriptActor:MonitorSpeed",
- ParentExecutionId = Guid.NewGuid(),
- };
- var original = new ReplicationOperation(ReplicationOperationType.Park, message.Id, message);
-
- var back = RoundTrip(original);
-
- Assert.Equal(original.OperationType, back.OperationType);
- Assert.Equal(original.MessageId, back.MessageId);
- Assert.NotNull(back.Message);
- var m = back.Message!;
- Assert.Equal(message.Id, m.Id);
- Assert.Equal(message.Category, m.Category);
- Assert.Equal(message.Target, m.Target);
- Assert.Equal(message.PayloadJson, m.PayloadJson);
- Assert.Equal(message.RetryCount, m.RetryCount);
- Assert.Equal(message.MaxRetries, m.MaxRetries);
- Assert.Equal(message.RetryIntervalMs, m.RetryIntervalMs);
- Assert.Equal(message.CreatedAt, m.CreatedAt);
- Assert.Equal(message.LastAttemptAt, m.LastAttemptAt);
- Assert.Equal(message.Status, m.Status);
- Assert.Equal(message.LastError, m.LastError);
- Assert.Equal(message.OriginInstanceName, m.OriginInstanceName);
- Assert.Equal(message.ExecutionId, m.ExecutionId);
- Assert.Equal(message.SourceScript, m.SourceScript);
- Assert.Equal(message.ParentExecutionId, m.ParentExecutionId);
- }
-
- // Type-identity pins: the reflective-JSON wire embeds CLR type manifests, so a
- // rename/move of either type silently breaks standby replication across a
- // rolling upgrade. If one fails, you are making a wire-breaking change.
- [Theory]
- [InlineData(typeof(ReplicationOperation), "ZB.MOM.WW.ScadaBridge.StoreAndForward.ReplicationOperation")]
- [InlineData(typeof(StoreAndForwardMessage), "ZB.MOM.WW.ScadaBridge.StoreAndForward.StoreAndForwardMessage")]
- public void ReplicationContract_TypeIdentity_IsPinned(Type type, string expectedFullName) =>
- Assert.Equal(expectedFullName, type.FullName);
-}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardReplicationTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardReplicationTests.cs
deleted file mode 100644
index 8080ce78..00000000
--- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardReplicationTests.cs
+++ /dev/null
@@ -1,266 +0,0 @@
-using Microsoft.Extensions.Logging.Abstractions;
-using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
-using ZB.MOM.WW.ScadaBridge.TestSupport;
-
-namespace ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests;
-
-///
-/// StoreAndForward-001: the active node must forward every buffer operation
-/// (add / remove / park) to the standby via the ReplicationService, so a
-/// failover does not lose the buffer.
-///
-public class StoreAndForwardReplicationTests : IAsyncLifetime, IDisposable
-{
- private readonly TestLocalDb _localDb;
- private readonly StoreAndForwardStorage _storage;
- private readonly StoreAndForwardService _service;
- private readonly List _replicated = new();
-
- public StoreAndForwardReplicationTests()
- {
- _localDb = TestLocalDb.CreateTemp("ReplTests");
-
- _storage = new StoreAndForwardStorage(_localDb.Db, NullLogger.Instance);
-
- var options = new StoreAndForwardOptions
- {
- DefaultRetryInterval = TimeSpan.Zero,
- DefaultMaxRetries = 1,
- RetryTimerInterval = TimeSpan.FromMinutes(10),
- ReplicationEnabled = true,
- };
-
- var replication = new ReplicationService(options, NullLogger.Instance);
- replication.SetReplicationHandler(op =>
- {
- lock (_replicated) _replicated.Add(op);
- return Task.CompletedTask;
- });
-
- _service = new StoreAndForwardService(
- _storage, options, NullLogger.Instance, replication);
- }
-
- public async Task InitializeAsync() => await _storage.InitializeAsync();
- public Task DisposeAsync() => Task.CompletedTask;
- public void Dispose()
- {
- var path = _localDb.Path;
- _localDb.Dispose();
- TestLocalDb.DeleteFiles(path);
- }
-
- /// Replication is fire-and-forget (Task.Run); poll until the expected ops arrive.
- private async Task> WaitForReplicationAsync(int count)
- {
- for (var i = 0; i < 100; i++)
- {
- lock (_replicated)
- if (_replicated.Count >= count) return _replicated.ToList();
- await Task.Delay(20);
- }
- lock (_replicated) return _replicated.ToList();
- }
-
- [Fact]
- public async Task BufferingAMessage_ReplicatesAnAddOperation()
- {
- // No handler registered → message is buffered → an Add is replicated.
- var result = await _service.EnqueueAsync(
- StoreAndForwardCategory.ExternalSystem, "api", """{}""");
- Assert.True(result.WasBuffered);
-
- var ops = await WaitForReplicationAsync(1);
- Assert.Contains(ops, o =>
- o.OperationType == ReplicationOperationType.Add && o.MessageId == result.MessageId);
- }
-
- [Fact]
- public async Task SuccessfulRetry_ReplicatesARemoveOperation()
- {
- var calls = 0;
- _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
- _ => ++calls == 1
- ? throw new HttpRequestException("transient")
- : Task.FromResult(true));
-
- var result = await _service.EnqueueAsync(
- StoreAndForwardCategory.ExternalSystem, "api", """{}""");
- await _service.RetryPendingMessagesAsync();
-
- var ops = await WaitForReplicationAsync(2);
- Assert.Contains(ops, o => o.OperationType == ReplicationOperationType.Add);
- Assert.Contains(ops, o =>
- o.OperationType == ReplicationOperationType.Remove && o.MessageId == result.MessageId);
- }
-
- [Fact]
- public async Task ParkedMessage_ReplicatesAParkOperation()
- {
- _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
- _ => throw new HttpRequestException("always fails"));
-
- var result = await _service.EnqueueAsync(
- StoreAndForwardCategory.ExternalSystem, "api", """{}""", maxRetries: 1);
- await _service.RetryPendingMessagesAsync();
-
- var ops = await WaitForReplicationAsync(2);
- Assert.Contains(ops, o =>
- o.OperationType == ReplicationOperationType.Park && o.MessageId == result.MessageId);
- }
-
- ///
- /// StoreAndForward-016: an operator discarding a parked message must replicate
- /// a Remove so the standby's copy is also deleted (otherwise the discarded
- /// message reappears in the parked list after a failover).
- ///
- [Fact]
- public async Task DiscardingAParkedMessage_ReplicatesARemoveOperation()
- {
- _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
- _ => throw new HttpRequestException("always fails"));
-
- var result = await _service.EnqueueAsync(
- StoreAndForwardCategory.ExternalSystem, "api", """{}""", maxRetries: 1);
- await _service.RetryPendingMessagesAsync(); // -> parked
- await WaitForReplicationAsync(2);
-
- var discarded = await _service.DiscardParkedMessageAsync(result.MessageId);
- Assert.True(discarded);
-
- var ops = await WaitForReplicationAsync(3);
- Assert.Contains(ops, o =>
- o.OperationType == ReplicationOperationType.Remove && o.MessageId == result.MessageId);
- }
-
- ///
- /// StoreAndForward-016: an operator retrying a parked message must replicate a
- /// Requeue so the standby's copy moves back to Pending (otherwise it stays
- /// Parked on the standby and the operator's retry is lost across a failover).
- ///
- [Fact]
- public async Task RetryingAParkedMessage_ReplicatesARequeueOperation()
- {
- _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem,
- _ => throw new HttpRequestException("always fails"));
-
- var result = await _service.EnqueueAsync(
- StoreAndForwardCategory.ExternalSystem, "api", """{}""", maxRetries: 1);
- await _service.RetryPendingMessagesAsync(); // -> parked
- await WaitForReplicationAsync(2);
-
- var retried = await _service.RetryParkedMessageAsync(result.MessageId);
- Assert.True(retried);
-
- var ops = await WaitForReplicationAsync(3);
- var requeue = ops.SingleOrDefault(o =>
- o.OperationType == ReplicationOperationType.Requeue && o.MessageId == result.MessageId);
- Assert.NotNull(requeue);
- Assert.NotNull(requeue!.Message);
- Assert.Equal(StoreAndForwardMessageStatus.Pending, requeue.Message!.Status);
- }
-
- ///
- /// StoreAndForward-016: the standby applies a Requeue by moving its row back to
- /// Pending with retry_count = 0, mirroring the active node's local state.
- ///
- [Fact]
- public async Task ApplyReplicatedOperation_Requeue_MovesStandbyRowBackToPending()
- {
- var replication = new ReplicationService(
- new StoreAndForwardOptions { ReplicationEnabled = true },
- NullLogger.Instance);
-
- var parked = new StoreAndForwardMessage
- {
- Id = "requeue1",
- Category = StoreAndForwardCategory.ExternalSystem,
- Target = "api",
- PayloadJson = "{}",
- RetryCount = 5,
- MaxRetries = 1,
- RetryIntervalMs = 0,
- CreatedAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Parked,
- };
- await _storage.EnqueueAsync(parked);
-
- var requeued = new StoreAndForwardMessage
- {
- Id = parked.Id,
- Category = parked.Category,
- Target = parked.Target,
- PayloadJson = parked.PayloadJson,
- RetryCount = 0,
- MaxRetries = parked.MaxRetries,
- RetryIntervalMs = parked.RetryIntervalMs,
- CreatedAt = parked.CreatedAt,
- Status = StoreAndForwardMessageStatus.Pending,
- };
- await replication.ApplyReplicatedOperationAsync(
- new ReplicationOperation(ReplicationOperationType.Requeue, parked.Id, requeued),
- _storage);
-
- var row = await _storage.GetMessageByIdAsync(parked.Id);
- Assert.NotNull(row);
- Assert.Equal(StoreAndForwardMessageStatus.Pending, row!.Status);
- Assert.Equal(0, row.RetryCount);
- }
-
- private static ReplicationService NewReplicationService() =>
- new(new StoreAndForwardOptions { ReplicationEnabled = true },
- NullLogger.Instance);
-
- private static StoreAndForwardMessage NewMessage(string id) => new()
- {
- Id = id,
- Category = StoreAndForwardCategory.ExternalSystem,
- Target = "api",
- PayloadJson = "{}",
- RetryCount = 0,
- MaxRetries = 1,
- RetryIntervalMs = 0,
- CreatedAt = DateTimeOffset.UtcNow,
- Status = StoreAndForwardMessageStatus.Pending,
- };
-
- ///
- /// A Park (or Requeue) whose original Add was lost — e.g. the Add's
- /// fire-and-forget replication dropped — must still materialise the row on the
- /// standby. The full message rides in the operation, so an upsert self-heals;
- /// a blind UPDATE would affect 0 rows and the row would be gone forever.
- ///
- [Fact]
- public async Task ApplyReplicatedPark_WhenAddWasLost_MaterializesTheParkedRow()
- {
- var service = NewReplicationService();
- var msg = NewMessage("lost-add-1"); // never Added on this (standby) storage
- msg.Status = StoreAndForwardMessageStatus.Parked;
-
- await service.ApplyReplicatedOperationAsync(
- new ReplicationOperation(ReplicationOperationType.Park, msg.Id, msg), _storage);
-
- var row = await _storage.GetMessageByIdAsync("lost-add-1");
- Assert.NotNull(row); // pre-fix: blind UPDATE affected 0 rows, row is gone forever
- Assert.Equal(StoreAndForwardMessageStatus.Parked, row!.Status);
- }
-
- ///
- /// A duplicate Add (e.g. after Task 21's anti-entropy resync re-issues an Add
- /// the standby already holds) must not violate the PK — the upsert applies
- /// newest-wins instead of throwing.
- ///
- [Fact]
- public async Task ApplyReplicatedAdd_Twice_IsIdempotent_NewestWins()
- {
- var service = NewReplicationService();
- var msg = NewMessage("dup-add");
- await service.ApplyReplicatedOperationAsync(
- new ReplicationOperation(ReplicationOperationType.Add, msg.Id, msg), _storage);
- msg.RetryCount = 3;
- await service.ApplyReplicatedOperationAsync( // pre-fix: SqliteException PK violation
- new ReplicationOperation(ReplicationOperationType.Add, msg.Id, msg), _storage);
-
- Assert.Equal(3, (await _storage.GetMessageByIdAsync("dup-add"))!.RetryCount);
- }
-}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardSiteEventTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardSiteEventTests.cs
index 52144ee7..b21acfd8 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardSiteEventTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardSiteEventTests.cs
@@ -54,7 +54,7 @@ public class StoreAndForwardSiteEventTests : IAsyncLifetime, IDisposable
_service = new StoreAndForwardService(
_storage, _options, NullLogger.Instance,
- replication: null, cachedCallObserver: null, siteId: "site-a",
+ cachedCallObserver: null, siteId: "site-a",
siteEventLogger: _siteLog);
}
diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardStorageTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardStorageTests.cs
index 668a3007..18ce27c2 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardStorageTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardStorageTests.cs
@@ -769,16 +769,15 @@ public class StoreAndForwardStorageTests : IAsyncLifetime, IDisposable
Assert.True(truncated); // a third row exists beyond the limit
}
- [Fact]
- public async Task ReplaceAll_SwapsTheEntireBuffer_Atomically()
- {
- await _storage.EnqueueAsync(NewMsg("stale"));
- await _storage.ReplaceAllAsync(new[] { NewMsg("fresh-1"), NewMsg("fresh-2") });
-
- Assert.Null(await _storage.GetMessageByIdAsync("stale"));
- Assert.NotNull(await _storage.GetMessageByIdAsync("fresh-1"));
- Assert.NotNull(await _storage.GetMessageByIdAsync("fresh-2"));
- }
+ // ReplaceAll_SwapsTheEntireBuffer_Atomically was DELETED with ReplaceAllAsync in
+ // LocalDb Phase 2, and deliberately not replaced. It asserted a destructive
+ // delete-all-then-insert-all, which was the standby's anti-entropy apply: a resync
+ // overwrote the standby's copy with the active node's snapshot. sf_messages is now a
+ // replicated table, so a mass DELETE would be CAPTURED and shipped to the peer — the
+ // method is not merely unused, it is unsafe to keep. LocalDb's own snapshot resync
+ // merges per row under last-writer-wins and never deletes, which is what
+ // LocalDbConfigConvergenceTests.ANodeWithNewerLocalRows_KeepsThem_WhenAPeerSnapshotArrives
+ // now covers.
// ── Task 23: oldest-parked-age health signal ──