diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs
index 4553c82c..2b8dda17 100644
--- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs
+++ b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs
@@ -394,7 +394,7 @@ public class StoreAndForwardService
});
_retryTimer = new Timer(
- _ => Volatile.Write(ref _sweepTask, RetryPendingMessagesAsync()),
+ _ => KickSweep(),
null,
_options.RetryTimerInterval,
_options.RetryTimerInterval);
@@ -658,17 +658,37 @@ public class StoreAndForwardService
///
/// Kicks a background retry sweep now (fire-and-forget). No-op before
/// (storage may be uninitialized and the timer, whose
- /// presence gates this, is not yet set). Overlap-safe: the sweep's
- /// CAS makes a concurrent kick a cheap no-op.
+ /// presence gates this, is not yet set). Overlap-safe: publishes into
+ /// only when the CAS is won — see .
///
public void TriggerSweep()
{
if (_retryTimer == null) return;
- Volatile.Write(ref _sweepTask, RetryPendingMessagesAsync());
+ KickSweep();
+ }
+
+ /// The current drain handle — test seam for the N3 clobber regression.
+ internal Task? CurrentSweepTaskForTest => Volatile.Read(ref _sweepTask);
+
+ ///
+ /// Starts a sweep IF none is in flight, publishing the new task into
+ /// only when this call wins the
+ /// CAS. A kick that loses the CAS returns without touching —
+ /// pre-fix it overwrote the drain handle with an instantly-completed no-op, so
+ /// proceeded with disposal under a still-running sweep
+ /// (review 02 round 2, N3) — defeating the drain exactly when a sweep outlives a tick.
+ ///
+ private void KickSweep()
+ {
+ if (Interlocked.CompareExchange(ref _retryInProgress, 1, 0) != 0)
+ return;
+ Volatile.Write(ref _sweepTask, RunSweepOwnedAsync());
}
///
/// Background retry sweep. Processes all pending messages that are due for retry.
+ /// Self-CASes for direct callers (existing tests); the entry CAS is skipped when
+ /// ownership is already held by .
///
/// A task representing the asynchronous retry sweep.
internal async Task RetryPendingMessagesAsync()
@@ -676,7 +696,16 @@ public class StoreAndForwardService
// Prevent overlapping retry sweeps
if (Interlocked.CompareExchange(ref _retryInProgress, 1, 0) != 0)
return;
+ await RunSweepOwnedAsync();
+ }
+ ///
+ /// The actual retry sweep body. The caller MUST already own the
+ /// flag (won the CAS); this method releases it in its
+ /// finally. Never call directly without holding ownership.
+ ///
+ private async Task RunSweepOwnedAsync()
+ {
try
{
var gate = _deliveryGate;
diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs
index 4facffa7..3c7a4c94 100644
--- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs
+++ b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs
@@ -821,4 +821,67 @@ public class StoreAndForwardServiceTests : IAsyncLifetime, IDisposable
}
finally { await service.StopAsync(); }
}
+
+ // ── R2 T8: _sweepTask clobber (N3) ──
+
+ [Fact]
+ public async Task TriggerSweep_WhileSweepInFlight_DoesNotClobberTheDrainHandle()
+ {
+ var service = CreateService(retryTimerInterval: TimeSpan.FromHours(1));
+ var entered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem, async _ =>
+ {
+ entered.TrySetResult();
+ await release.Task;
+ return true;
+ });
+ await service.StartAsync();
+ try
+ {
+ await service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "t", "{}",
+ attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero);
+
+ service.TriggerSweep(); // real sweep, blocked in the handler
+ await entered.Task.WaitAsync(TimeSpan.FromSeconds(5));
+ service.TriggerSweep(); // redundant kick — pre-fix clobbers _sweepTask
+
+ var handle = service.CurrentSweepTaskForTest;
+ Assert.NotNull(handle);
+ Assert.False(handle!.IsCompleted); // pre-fix: true (a completed no-op replaced the real sweep)
+ }
+ finally
+ {
+ release.TrySetResult();
+ await service.StopAsync();
+ }
+ }
+
+ [Fact]
+ public async Task StopAsync_WaitsForTheRealInFlightSweep_EvenAfterARedundantTrigger()
+ {
+ var service = CreateService(retryTimerInterval: TimeSpan.FromHours(1));
+ var entered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem, async _ =>
+ {
+ entered.TrySetResult();
+ await release.Task;
+ return true;
+ });
+ await service.StartAsync();
+ await service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "t", "{}",
+ attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero);
+
+ service.TriggerSweep();
+ await entered.Task.WaitAsync(TimeSpan.FromSeconds(5));
+ service.TriggerSweep(); // the clobbering kick
+
+ var stop = service.StopAsync();
+ await Task.Delay(300);
+ Assert.False(stop.IsCompleted); // pre-fix: StopAsync already returned (awaited the no-op)
+
+ release.TrySetResult();
+ await stop.WaitAsync(TimeSpan.FromSeconds(5)); // drains the real sweep promptly once released
+ }
}