diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs index d2aad773..77527dd5 100644 --- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs +++ b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs @@ -85,6 +85,21 @@ public class StoreAndForwardService private Timer? _retryTimer; private int _retryInProgress; + /// + /// Active-node delivery gate. When set, the retry sweep runs ONLY when the + /// gate reports true — the standby node applies replicated buffer operations + /// but must never deliver (Component-StoreAndForward.md "standby-passive" + /// model). Re-evaluated on every sweep tick so a failover resumes delivery + /// within one RetryTimerInterval without re-wiring. Null (tests, central + /// hosts) preserves the ungated legacy behaviour. A throwing gate is treated + /// as standby — safe-by-default, mirroring SiteCommunicationActor's + /// DefaultIsActiveCheck fallback. + /// + private Func? _deliveryGate; + + /// Installs the active-node delivery gate (see ). + public void SetDeliveryGate(Func gate) => _deliveryGate = gate; + /// /// The in-flight retry sweep , or /// null when no sweep is currently running. Captured when the timer @@ -568,6 +583,24 @@ public class StoreAndForwardService try { + var gate = _deliveryGate; + if (gate != null) + { + bool isActive; + try { isActive = gate(); } + catch (Exception ex) + { + _logger.LogWarning(ex, + "S&F delivery gate threw; treating this node as standby for this sweep"); + isActive = false; + } + if (!isActive) + { + _logger.LogDebug("S&F retry sweep skipped: this node is not the active site node"); + return; + } + } + var messages = await _storage.GetMessagesForRetryAsync(); if (messages.Count == 0) return; diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs index 9d479afc..24454a27 100644 --- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs +++ b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs @@ -582,4 +582,60 @@ public class StoreAndForwardServiceTests : IAsyncLifetime, IDisposable Assert.True(handlerCompleted, "Sweep handler must have finished before StopAsync returned."); } + + // ── Task 3 (arch review 02, Stability #2): active-node delivery gate ── + // The standby site node applies replicated buffer operations but must never + // deliver; the sweep runs only when the gate reports the node is active. + + [Fact] + public async Task RetrySweep_SkipsDelivery_WhenDeliveryGateReportsStandby() + { + var delivered = 0; + _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem, + _ => { delivered++; return Task.FromResult(true); }); + _service.SetDeliveryGate(() => false); // standby + + var result = await _service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "target-x", "{}", + attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero); + await _service.RetryPendingMessagesAsync(); + + Assert.Equal(0, delivered); + var row = await _service.GetMessageByIdAsync(result.MessageId); + Assert.NotNull(row); + Assert.Equal(StoreAndForwardMessageStatus.Pending, row!.Status); // row untouched + } + + [Fact] + public async Task RetrySweep_ResumesDelivery_WhenGateFlipsActive() + { + var delivered = 0; + var active = false; + _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem, + _ => { delivered++; return Task.FromResult(true); }); + _service.SetDeliveryGate(() => active); + + await _service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "target-x", "{}", + attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero); + + await _service.RetryPendingMessagesAsync(); + Assert.Equal(0, delivered); + + active = true; // failover: this node became active + await _service.RetryPendingMessagesAsync(); + Assert.Equal(1, delivered); + } + + [Fact] + public async Task RetrySweep_TreatsThrowingGateAsStandby() + { + var delivered = 0; + _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem, + _ => { delivered++; return Task.FromResult(true); }); + _service.SetDeliveryGate(() => throw new InvalidOperationException("cluster not ready")); + + await _service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "t", "{}", + attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero); + await _service.RetryPendingMessagesAsync(); // must not throw + Assert.Equal(0, delivered); + } }