diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardOptions.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardOptions.cs index 7546e942..39aaa321 100644 --- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardOptions.cs +++ b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardOptions.cs @@ -29,4 +29,10 @@ public class StoreAndForwardOptions /// 0 = unlimited (legacy). /// public int SweepBatchLimit { get; set; } = 500; + + /// + /// Max (category, target) lanes delivered concurrently per sweep. + /// 1 = legacy serial. Within a lane delivery stays sequential (per-target FIFO). + /// + public int SweepTargetParallelism { get; set; } = 4; } diff --git a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs index d3ac792d..2f1ef7b6 100644 --- a/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs +++ b/src/ZB.MOM.WW.ScadaBridge.StoreAndForward/StoreAndForwardService.cs @@ -606,20 +606,36 @@ public class StoreAndForwardService _logger.LogDebug("Retry sweep: {Count} messages due for retry", messages.Count); - var failedTargets = new HashSet<(StoreAndForwardCategory, string)>(); - foreach (var message in messages) - { - // One transient failure per (category, target) per sweep: the target - // is down — burning a full timeout per remaining message serializes - // the sweep into hours under backlog (arch review 02, Performance #1). - // Skipped rows keep their RetryCount/LastAttemptAt untouched. - if (failedTargets.Contains((message.Category, message.Target))) - continue; + // Group due rows into per-(category,target) lanes. GroupBy is stable, so + // each lane preserves created_at-ASC (oldest-first) order. Lanes run + // concurrently up to SweepTargetParallelism, so one slow/dead target no + // longer blocks delivery to healthy ones (arch review 02, Performance #2); + // within a lane delivery stays strictly sequential (per-target FIFO). + var lanes = messages + .GroupBy(m => (m.Category, m.Target)) + .Select(g => g.ToList()) + .ToList(); - var outcome = await RetryMessageAsync(message); - if (outcome == RetryOutcome.TransientFailure) - failedTargets.Add((message.Category, message.Target)); - } + using var laneCap = new SemaphoreSlim(Math.Max(1, _options.SweepTargetParallelism)); + var laneTasks = lanes.Select(async lane => + { + await laneCap.WaitAsync(); + try + { + foreach (var message in lane) + { + // Task 8 short-circuit, per lane: one transient failure means + // the target is down — skip the rest of this lane this sweep + // rather than burning a full timeout per remaining message. + // Skipped rows keep their RetryCount/LastAttemptAt untouched. + var outcome = await RetryMessageAsync(message); + if (outcome == RetryOutcome.TransientFailure) + return; + } + } + finally { laneCap.Release(); } + }).ToList(); + await Task.WhenAll(laneTasks); } catch (Exception ex) { diff --git a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs index 7ffcaefe..8fd5fda9 100644 --- a/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs +++ b/tests/ZB.MOM.WW.ScadaBridge.StoreAndForward.Tests/StoreAndForwardServiceTests.cs @@ -686,4 +686,64 @@ public class StoreAndForwardServiceTests : IAsyncLifetime, IDisposable // assertion doesn't depend on created_at tie-breaking. Assert.Equal(new[] { 0, 1 }, new[] { r1, r2 }.OrderBy(x => x).ToArray()); } + + // ── Task 9 (arch review 02, Performance): parallel per-target lanes ── + // A slow/dead target no longer serializes the whole sweep; delivery within a + // single (category, target) lane stays sequential (per-target FIFO). + + [Fact] + public async Task RetrySweep_SlowTarget_DoesNotBlockOtherTargets() + { + var slowGate = new TaskCompletionSource(); + var healthyDelivered = new TaskCompletionSource(); + _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem, async msg => + { + if (msg.Target == "slow") { await slowGate.Task; return true; } + healthyDelivered.TrySetResult(); + return true; + }); + await _service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "slow", "{}", + attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero); + await _service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "healthy", "{}", + attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero); + + var sweep = _service.RetryPendingMessagesAsync(); + // Healthy lane completes while the slow lane is still blocked — pre-fix this + // times out because delivery is strictly serial across targets. + await healthyDelivered.Task.WaitAsync(TimeSpan.FromSeconds(5)); + slowGate.SetResult(); + await sweep; + } + + [Fact] + public async Task RetrySweep_WithinTargetLane_StaysSequential() + { + var concurrent = 0; + var maxConcurrent = 0; + var delivered = 0; + _service.RegisterDeliveryHandler(StoreAndForwardCategory.ExternalSystem, async _ => + { + var c = Interlocked.Increment(ref concurrent); + InterlockedMax(ref maxConcurrent, c); + await Task.Delay(20); + Interlocked.Decrement(ref concurrent); + Interlocked.Increment(ref delivered); + return true; + }); + for (var i = 0; i < 4; i++) + await _service.EnqueueAsync(StoreAndForwardCategory.ExternalSystem, "same-target", "{}", + attemptImmediateDelivery: false, retryInterval: TimeSpan.Zero, messageId: $"m{i}"); + + await _service.RetryPendingMessagesAsync(); + + Assert.Equal(4, delivered); + Assert.Equal(1, maxConcurrent); // one (category,target) lane: strictly sequential + } + + private static void InterlockedMax(ref int target, int value) + { + int current; + do { current = Volatile.Read(ref target); if (value <= current) return; } + while (Interlocked.CompareExchange(ref target, value, current) != current); + } }