fix(store-and-forward): gate the retry sweep behind an active-node delivery gate (standby must be passive)

This commit is contained in:
Joseph Doherty
2026-07-08 17:30:21 -04:00
parent ec47cb5612
commit 76c10c5de6
2 changed files with 89 additions and 0 deletions
@@ -85,6 +85,21 @@ public class StoreAndForwardService
private Timer? _retryTimer; private Timer? _retryTimer;
private int _retryInProgress; private int _retryInProgress;
/// <summary>
/// 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.
/// </summary>
private Func<bool>? _deliveryGate;
/// <summary>Installs the active-node delivery gate (see <see cref="_deliveryGate"/>).</summary>
public void SetDeliveryGate(Func<bool> gate) => _deliveryGate = gate;
/// <summary> /// <summary>
/// The in-flight retry sweep <see cref="Task"/>, or /// The in-flight retry sweep <see cref="Task"/>, or
/// <c>null</c> when no sweep is currently running. Captured when the timer /// <c>null</c> when no sweep is currently running. Captured when the timer
@@ -568,6 +583,24 @@ public class StoreAndForwardService
try 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(); var messages = await _storage.GetMessagesForRetryAsync();
if (messages.Count == 0) return; if (messages.Count == 0) return;
@@ -582,4 +582,60 @@ public class StoreAndForwardServiceTests : IAsyncLifetime, IDisposable
Assert.True(handlerCompleted, Assert.True(handlerCompleted,
"Sweep handler must have finished before StopAsync returned."); "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);
}
} }