perf(store-and-forward): parallel per-target sweep lanes (cap 4), sequential within lane
This commit is contained in:
@@ -29,4 +29,10 @@ public class StoreAndForwardOptions
|
|||||||
/// 0 = unlimited (legacy).
|
/// 0 = unlimited (legacy).
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public int SweepBatchLimit { get; set; } = 500;
|
public int SweepBatchLimit { get; set; } = 500;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Max <c>(category, target)</c> lanes delivered concurrently per sweep.
|
||||||
|
/// 1 = legacy serial. Within a lane delivery stays sequential (per-target FIFO).
|
||||||
|
/// </summary>
|
||||||
|
public int SweepTargetParallelism { get; set; } = 4;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -606,20 +606,36 @@ public class StoreAndForwardService
|
|||||||
|
|
||||||
_logger.LogDebug("Retry sweep: {Count} messages due for retry", messages.Count);
|
_logger.LogDebug("Retry sweep: {Count} messages due for retry", messages.Count);
|
||||||
|
|
||||||
var failedTargets = new HashSet<(StoreAndForwardCategory, string)>();
|
// Group due rows into per-(category,target) lanes. GroupBy is stable, so
|
||||||
foreach (var message in messages)
|
// each lane preserves created_at-ASC (oldest-first) order. Lanes run
|
||||||
{
|
// concurrently up to SweepTargetParallelism, so one slow/dead target no
|
||||||
// One transient failure per (category, target) per sweep: the target
|
// longer blocks delivery to healthy ones (arch review 02, Performance #2);
|
||||||
// is down — burning a full timeout per remaining message serializes
|
// within a lane delivery stays strictly sequential (per-target FIFO).
|
||||||
// the sweep into hours under backlog (arch review 02, Performance #1).
|
var lanes = messages
|
||||||
// Skipped rows keep their RetryCount/LastAttemptAt untouched.
|
.GroupBy(m => (m.Category, m.Target))
|
||||||
if (failedTargets.Contains((message.Category, message.Target)))
|
.Select(g => g.ToList())
|
||||||
continue;
|
.ToList();
|
||||||
|
|
||||||
var outcome = await RetryMessageAsync(message);
|
using var laneCap = new SemaphoreSlim(Math.Max(1, _options.SweepTargetParallelism));
|
||||||
if (outcome == RetryOutcome.TransientFailure)
|
var laneTasks = lanes.Select(async lane =>
|
||||||
failedTargets.Add((message.Category, message.Target));
|
{
|
||||||
}
|
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)
|
catch (Exception ex)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -686,4 +686,64 @@ public class StoreAndForwardServiceTests : IAsyncLifetime, IDisposable
|
|||||||
// assertion doesn't depend on created_at tie-breaking.
|
// assertion doesn't depend on created_at tie-breaking.
|
||||||
Assert.Equal(new[] { 0, 1 }, new[] { r1, r2 }.OrderBy(x => x).ToArray());
|
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);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user