diff --git a/framework/src/BBT.Aether.Core/BBT/Aether/BackgroundJob/IBackgroundJobArmHandle.cs b/framework/src/BBT.Aether.Core/BBT/Aether/BackgroundJob/IBackgroundJobArmHandle.cs
new file mode 100644
index 0000000..6a0c37b
--- /dev/null
+++ b/framework/src/BBT.Aether.Core/BBT/Aether/BackgroundJob/IBackgroundJobArmHandle.cs
@@ -0,0 +1,33 @@
+using System;
+using System.Threading;
+using System.Threading.Tasks;
+
+namespace BBT.Aether.BackgroundJob;
+
+///
+/// A deferred arm for an already-persisted background job: everything the external scheduler needs is
+/// captured in memory, so arming later costs exactly one scheduler call — no re-read of the job row and
+/// no extra status write.
+///
+/// Exists for callers that must persist the job inside a critical section but cannot afford to make the
+/// scheduler round-trip there. Holding a distributed lock across an external call makes that call the
+/// lock's hold time, serializing every other contender behind it.
+///
+///
+/// The row is already Scheduled when the handle is issued — optimistically, because the common
+/// case succeeds. reconciles a failure by rolling the row back to
+/// Pending so the arming poller reclaims it, which is the same contract the inline arm has.
+///
+///
+public interface IBackgroundJobArmHandle
+{
+ /// The id of the persisted job this handle arms.
+ Guid JobId { get; }
+
+ ///
+ /// Arms the job in the external scheduler. Never throws: a failure is logged and the row is rolled
+ /// back to Pending for the arming poller. Safe to call once; calling it again re-schedules
+ /// the same job name, which the scheduler treats as an overwrite.
+ ///
+ Task ArmAsync(CancellationToken cancellationToken = default);
+}
diff --git a/framework/src/BBT.Aether.Core/BBT/Aether/BackgroundJob/IBackgroundJobService.cs b/framework/src/BBT.Aether.Core/BBT/Aether/BackgroundJob/IBackgroundJobService.cs
index 3210358..c54fcc2 100644
--- a/framework/src/BBT.Aether.Core/BBT/Aether/BackgroundJob/IBackgroundJobService.cs
+++ b/framework/src/BBT.Aether.Core/BBT/Aether/BackgroundJob/IBackgroundJobService.cs
@@ -101,4 +101,33 @@ Task CancelWaitingAsync(
///
/// Thrown when id is empty.
Task DeleteAsync(Guid id, CancellationToken cancellationToken = default);
+
+ ///
+ /// Persists the job and returns a handle that arms it later, instead of arming as part of this call.
+ ///
+ /// For callers that write the job inside a critical section (a distributed lock, a short
+ /// transaction) but must keep the scheduler round-trip out of it. Persisting is cheap and local;
+ /// the scheduler call is neither, and inside a lock it becomes the lock's hold time.
+ ///
+ ///
+ /// Same parameters and same persistence semantics as with
+ /// directly: true — the row lands Scheduled and an arm failure rolls it back to
+ /// Pending for the arming poller. The only difference is WHEN the scheduler is called, which
+ /// the caller now decides by invoking .
+ ///
+ ///
+ /// Call ArmAsync only after the work that justified the critical section has committed. The
+ /// handle carries the payload in memory, so arming costs one scheduler call and no database access.
+ ///
+ ///
+ Task EnqueueWithDeferredArmAsync(
+ string handlerName,
+ string jobName,
+ TPayload payload,
+ string schedule,
+ Dictionary? metadata = null,
+ JobScheduleFailurePolicy? failurePolicyOptions = null,
+ Guid? jobId = null,
+ BBT.Aether.Domain.Entities.JobKind? kind = null,
+ CancellationToken cancellationToken = default);
}
diff --git a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/BackgroundJobService.cs b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/BackgroundJobService.cs
index e5a830b..8cf8199 100644
--- a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/BackgroundJobService.cs
+++ b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/BackgroundJobService.cs
@@ -78,6 +78,46 @@ public async Task EnqueueAsync(
Guid? jobId = null,
JobKind? kind = null,
CancellationToken cancellationToken = default)
+ {
+ var (enqueuedId, _) = await EnqueueCoreAsync(
+ handlerName, jobName, payload, schedule, metadata, failurePolicyOptions,
+ directly, jobId, kind, deferArm: false, cancellationToken);
+ return enqueuedId;
+ }
+
+ ///
+ public async Task EnqueueWithDeferredArmAsync(
+ string handlerName,
+ string jobName,
+ TPayload payload,
+ string schedule,
+ Dictionary? metadata = null,
+ JobScheduleFailurePolicy? failurePolicyOptions = null,
+ Guid? jobId = null,
+ JobKind? kind = null,
+ CancellationToken cancellationToken = default)
+ {
+ // directly: true so the row lands Scheduled, exactly as the inline path leaves it — no later
+ // status write is needed. The one difference is that the scheduler is not called here; the
+ // returned handle calls it, carrying the payload in memory so no re-read is needed either.
+ var (deferredId, arm) = await EnqueueCoreAsync(
+ handlerName, jobName, payload, schedule, metadata, failurePolicyOptions,
+ directly: true, jobId, kind, deferArm: true, cancellationToken);
+ return new DeferredArmHandle(deferredId, arm!);
+ }
+
+ private async Task<(Guid JobId, Func? Arm)> EnqueueCoreAsync(
+ string handlerName,
+ string jobName,
+ TPayload payload,
+ string schedule,
+ Dictionary? metadata,
+ JobScheduleFailurePolicy? failurePolicyOptions,
+ bool directly,
+ Guid? jobId,
+ JobKind? kind,
+ bool deferArm,
+ CancellationToken cancellationToken)
{
if (string.IsNullOrWhiteSpace(handlerName))
throw new ArgumentNullException(nameof(handlerName));
@@ -91,7 +131,10 @@ public async Task EnqueueAsync(
if (string.IsNullOrWhiteSpace(schedule))
throw new ArgumentNullException(nameof(schedule));
- using var activity = InfrastructureActivitySource.Source.StartActivity(
+ // Diagnostic, not business: enqueueing a job is plumbing around the work, and in the
+ // default Business profile it only adds depth to an already deep trace. The job's own
+ // execution span (BackgroundJob.Execute) is the one that carries meaning and stays.
+ using var activity = InfrastructureActivitySource.StartDiagnosticActivity(
"BackgroundJob.Enqueue",
ActivityKind.Producer,
Activity.Current?.Context ?? default);
@@ -155,19 +198,26 @@ public async Task EnqueueAsync(
// Bytes for the scheduler (the `directly` arm path). Equivalent to the JSON the poller arms with.
var payloadBytes = eventSerializer.Serialize(envelope);
+ // Deferred arm: hand the caller a closure over everything the scheduler needs. Same call the
+ // inline path makes, same failure handling — only the timing is the caller's to choose.
+ Func? armAction = deferArm
+ ? ct => ArmNowAsync(handlerName, jobName, schedule, payloadBytes,
+ failurePolicyOptions, effectiveJobId, ct)
+ : null;
+
// Atomic-ambient: when the caller has an ambient UoW, persist into it (commits with their business
// transaction — a rollback discards the row). Otherwise open a short own transaction.
if (uowManager.Current is { } ambient)
{
await jobStore.SaveAsync(jobInfo, cancellationToken);
- if (directly)
+ if (directly && !deferArm)
ambient.OnCompleted(_ => ArmNowAsync(handlerName, jobName, schedule, payloadBytes,
failurePolicyOptions, effectiveJobId, CancellationToken.None));
logger.LogInformation(
"Enqueued {Status} job '{HandlerName}'/'{JobName}' into ambient UoW. Id: {Id}",
jobInfo.Status, handlerName, jobName, effectiveJobId);
activity?.SetStatus(ActivityStatusCode.Ok);
- return effectiveJobId;
+ return (effectiveJobId, armAction);
}
await using (var uow = uowManager.Begin(
@@ -187,14 +237,26 @@ public async Task EnqueueAsync(
}
}
- if (directly)
+ if (directly && !deferArm)
await ArmNowAsync(handlerName, jobName, schedule, payloadBytes, failurePolicyOptions,
effectiveJobId, cancellationToken);
logger.LogInformation(
"Enqueued {Status} job '{HandlerName}'/'{JobName}'. Id: {Id}",
jobInfo.Status, handlerName, jobName, effectiveJobId);
activity?.SetStatus(ActivityStatusCode.Ok);
- return effectiveJobId;
+ return (effectiveJobId, armAction);
+ }
+
+ ///
+ /// Closure-backed . Holds the scheduler arguments captured at
+ /// enqueue time, so arming needs neither a job-row read nor a status write.
+ ///
+ private sealed class DeferredArmHandle(Guid jobId, Func arm)
+ : IBackgroundJobArmHandle
+ {
+ public Guid JobId { get; } = jobId;
+
+ public Task ArmAsync(CancellationToken cancellationToken = default) => arm(cancellationToken);
}
///
@@ -248,7 +310,10 @@ public async Task UpdateAsync(Guid id, string newSchedule, CancellationToken can
if (string.IsNullOrWhiteSpace(newSchedule))
throw new ArgumentNullException(nameof(newSchedule));
- using var activity = InfrastructureActivitySource.Source.StartActivity(
+ // Diagnostic, not business: enqueueing a job is plumbing around the work, and in the
+ // default Business profile it only adds depth to an already deep trace. The job's own
+ // execution span (BackgroundJob.Execute) is the one that carries meaning and stays.
+ using var activity = InfrastructureActivitySource.StartDiagnosticActivity(
"BackgroundJob.Update",
ActivityKind.Producer,
Activity.Current?.Context ?? default);
@@ -413,7 +478,10 @@ public async Task DeleteAsync(Guid id, CancellationToken cancellationToken
if (id == Guid.Empty)
throw new ArgumentException("Id cannot be empty.", nameof(id));
- using var activity = InfrastructureActivitySource.Source.StartActivity(
+ // Diagnostic, not business: enqueueing a job is plumbing around the work, and in the
+ // default Business profile it only adds depth to an already deep trace. The job's own
+ // execution span (BackgroundJob.Execute) is the one that carries meaning and stays.
+ using var activity = InfrastructureActivitySource.StartDiagnosticActivity(
"BackgroundJob.Delete",
ActivityKind.Producer,
Activity.Current?.Context ?? default);
diff --git a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Dapr/DaprJobScheduler.cs b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Dapr/DaprJobScheduler.cs
index bc73481..521a6df 100644
--- a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Dapr/DaprJobScheduler.cs
+++ b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Dapr/DaprJobScheduler.cs
@@ -146,7 +146,9 @@ await daprJobsClient.DeleteJobAsync(
private static Activity? StartSchedulerActivity(string operationName, string handlerName, string jobName)
{
- var activity = InfrastructureActivitySource.Source.StartActivity(
+ // Diagnostic: the scheduler round-trip is infrastructure detail. Business traces care that
+ // the job ran (BackgroundJob.Execute), not about the Schedule/Delete calls that armed it.
+ var activity = InfrastructureActivitySource.StartDiagnosticActivity(
operationName,
ActivityKind.Client,
Activity.Current?.Context ?? default);
diff --git a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/JobDispatcher.cs b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/JobDispatcher.cs
index fbf27bf..24c3811 100644
--- a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/JobDispatcher.cs
+++ b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/JobDispatcher.cs
@@ -32,7 +32,8 @@ public virtual async Task DispatchAsync(
if (string.IsNullOrWhiteSpace(jobName))
throw new ArgumentNullException(nameof(jobName));
- using var activity = InfrastructureActivitySource.Source.StartActivity(
+ // Diagnostic: pure dispatch plumbing between BackgroundJob.Execute and the handler.
+ using var activity = InfrastructureActivitySource.StartDiagnosticActivity(
"BackgroundJob.Dispatch",
ActivityKind.Internal,
Activity.Current?.Context ?? default);
diff --git a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Processing/BackgroundJobArmingHostedService.cs b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Processing/BackgroundJobArmingHostedService.cs
index 46479fd..5fc24fb 100644
--- a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Processing/BackgroundJobArmingHostedService.cs
+++ b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/BackgroundJob/Processing/BackgroundJobArmingHostedService.cs
@@ -1,6 +1,7 @@
using System;
using System.Threading;
using System.Threading.Tasks;
+using BBT.Aether.Polling;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
@@ -9,8 +10,15 @@ namespace BBT.Aether.BackgroundJob.Processing;
///
/// Hosted service that drives on a timer. Each tick runs one
/// arming pass; exceptions per tick are caught and logged so a transient failure never tears down the
-/// loop. The delay between ticks is . Registered by the
-/// DI wiring (see AddAetherBackgroundJob); not auto-registered here.
+/// loop. The delay between ticks is , jittered.
+/// Registered by the DI wiring (see AddAetherBackgroundJob); not auto-registered here.
+///
+/// The interval is fixed — there is no adaptive backoff here, because an unarmed job must be picked up
+/// within a bounded time regardless of how quiet the system is. That makes jitter the only thing
+/// keeping replicas apart: without it, pods started together by a rolling deployment run every pass in
+/// lockstep, turning each tick into a burst of simultaneous claim queries over the same rows. A random
+/// startup offset spreads the first pass as well.
+///
///
public class BackgroundJobArmingHostedService(
BackgroundJobArmingProcessor processor,
@@ -24,6 +32,15 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
"Background-job arming poller started (interval {Interval}, schema {Schema}).",
options.ArmingInterval, options.Schema);
+ try
+ {
+ await Task.Delay(PollingDelay.StartupOffset(options.ArmingInterval), stoppingToken);
+ }
+ catch (OperationCanceledException)
+ {
+ return;
+ }
+
while (!stoppingToken.IsCancellationRequested)
{
try
@@ -41,7 +58,7 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
try
{
- await Task.Delay(options.ArmingInterval, stoppingToken);
+ await Task.Delay(PollingDelay.Jitter(options.ArmingInterval), stoppingToken);
}
catch (OperationCanceledException)
{
diff --git a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/InboxBackgroundService.cs b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/InboxBackgroundService.cs
index f7be54e..69ad397 100644
--- a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/InboxBackgroundService.cs
+++ b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/InboxBackgroundService.cs
@@ -1,6 +1,7 @@
using System;
using System.Threading;
using System.Threading.Tasks;
+using BBT.Aether.Polling;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
@@ -13,6 +14,18 @@ public sealed class InboxBackgroundService(
{
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
+ // Spread the first poll of replicas that booted together, so a rolling deployment does not
+ // leave the whole fleet polling on the same tick.
+ try
+ {
+ await Task.Delay(PollingDelay.StartupOffset(options.IdlePollingInterval), stoppingToken)
+ .ConfigureAwait(false);
+ }
+ catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
+ {
+ return;
+ }
+
var delay = options.IdlePollingInterval;
while (!stoppingToken.IsCancellationRequested)
@@ -21,8 +34,8 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
var processed = await processor.RunAsync(stoppingToken);
delay = processed > 0
- ? options.BusyPollingInterval
- : Min(delay * 2, options.MaxPollingInterval);
+ ? PollingDelay.OnProcessed(options.BusyPollingInterval)
+ : PollingDelay.OnEmpty(delay, options.MaxPollingInterval);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
@@ -31,12 +44,12 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
catch (Exception ex)
{
logger.LogError(ex, "Inbox background service error");
- delay = options.MaxPollingInterval;
+ // One step back, not straight to the cap: a transient fault must not stall every
+ // replica for a full maximum interval.
+ delay = PollingDelay.OnError(delay, options.IdlePollingInterval, options.MaxPollingInterval);
}
- await Task.Delay(delay, stoppingToken).ConfigureAwait(false);
+ await Task.Delay(PollingDelay.Jitter(delay), stoppingToken).ConfigureAwait(false);
}
}
-
- private static TimeSpan Min(TimeSpan a, TimeSpan b) => a < b ? a : b;
}
diff --git a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/OutboxBackgroundService.cs b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/OutboxBackgroundService.cs
index bf71cf7..8a92b90 100644
--- a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/OutboxBackgroundService.cs
+++ b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Events/Processing/OutboxBackgroundService.cs
@@ -1,6 +1,7 @@
using System;
using System.Threading;
using System.Threading.Tasks;
+using BBT.Aether.Polling;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
@@ -13,6 +14,18 @@ public sealed class OutboxBackgroundService(
{
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
+ // Spread the first poll of replicas that booted together, so a rolling deployment does not
+ // leave the whole fleet polling on the same tick.
+ try
+ {
+ await Task.Delay(PollingDelay.StartupOffset(options.IdlePollingInterval), stoppingToken)
+ .ConfigureAwait(false);
+ }
+ catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
+ {
+ return;
+ }
+
var delay = options.IdlePollingInterval;
while (!stoppingToken.IsCancellationRequested)
@@ -21,8 +34,8 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
var processed = await processor.RunAsync(stoppingToken);
delay = processed > 0
- ? options.BusyPollingInterval
- : Min(delay * 2, options.MaxPollingInterval);
+ ? PollingDelay.OnProcessed(options.BusyPollingInterval)
+ : PollingDelay.OnEmpty(delay, options.MaxPollingInterval);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
@@ -31,12 +44,12 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
catch (Exception ex)
{
logger.LogError(ex, "Outbox background service error");
- delay = options.MaxPollingInterval;
+ // One step back, not straight to the cap: a transient fault must not stall every
+ // replica for a full maximum interval.
+ delay = PollingDelay.OnError(delay, options.IdlePollingInterval, options.MaxPollingInterval);
}
- await Task.Delay(delay, stoppingToken).ConfigureAwait(false);
+ await Task.Delay(PollingDelay.Jitter(delay), stoppingToken).ConfigureAwait(false);
}
}
-
- private static TimeSpan Min(TimeSpan a, TimeSpan b) => a < b ? a : b;
}
diff --git a/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Polling/PollingDelay.cs b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Polling/PollingDelay.cs
new file mode 100644
index 0000000..5b1b05f
--- /dev/null
+++ b/framework/src/BBT.Aether.Infrastructure/BBT/Aether/Polling/PollingDelay.cs
@@ -0,0 +1,77 @@
+using System;
+
+namespace BBT.Aether.Polling;
+
+///
+/// Adaptive poll pacing shared by the outbox, inbox and background-job arming loops. Lives here rather than inline in each
+/// loop so the two stay identical and so the pacing rules are unit-testable without a running host.
+///
+///
+/// Every returned delay is jittered. Without it, replicas started together — which is exactly what a
+/// rolling deployment produces — poll in lockstep: the fleet loses the natural staggering that makes
+/// N replicas pick work up ~N times sooner than one, and each tick becomes a burst of simultaneous
+/// claim queries against the same rows. Jitter is what keeps the phases spread.
+///
+internal static class PollingDelay
+{
+ /// Jitter applied to every delay, as a fraction either side of the nominal value.
+ internal const double JitterFraction = 0.25;
+
+ private static readonly TimeSpan Floor = TimeSpan.FromMilliseconds(1);
+
+ ///
+ /// The delay after a round that found work: poll again almost immediately, since a queue that had
+ /// one item usually has more.
+ ///
+ internal static TimeSpan OnProcessed(TimeSpan busyInterval) => busyInterval;
+
+ ///
+ /// The delay after an empty round: double it, capped, so a quiet system stops paying for polls.
+ ///
+ internal static TimeSpan OnEmpty(TimeSpan current, TimeSpan max) => MinOf(Double(current), max);
+
+ ///
+ /// The delay after a failed round. Backs off one step like an empty round, but never below
+ /// so a hard failure is not retried at the busy cadence.
+ ///
+ /// Deliberately NOT a jump straight to . That is what the loops used to do,
+ /// and with several replicas a single transient fault — a brief database hiccup, one poison
+ /// message — stalled the entire fleet for a full maximum interval, right when it was most needed.
+ /// Escalating instead keeps a one-off blip cheap while a persistent fault still ends up at the cap.
+ ///
+ ///
+ internal static TimeSpan OnError(TimeSpan current, TimeSpan idleInterval, TimeSpan max)
+ => MinOf(MaxOf(Double(current), idleInterval), max);
+
+ ///
+ /// Applies to a delay using a caller-supplied uniform sample in
+ /// [0, 1), so tests can pin the arithmetic.
+ ///
+ internal static TimeSpan Jitter(TimeSpan delay, double sample)
+ {
+ var scale = 1.0 - JitterFraction + (2.0 * JitterFraction * sample);
+ var jittered = TimeSpan.FromTicks((long)(delay.Ticks * scale));
+ return jittered < Floor ? Floor : jittered;
+ }
+
+ /// Applies jitter using the shared random source.
+ internal static TimeSpan Jitter(TimeSpan delay) => Jitter(delay, Random.Shared.NextDouble());
+
+ ///
+ /// A random delay in [0, ) to spread the first poll of replicas that
+ /// started at the same moment. Bounded by the idle interval, not the maximum, so a fresh
+ /// deployment never sits idle for a whole cap before its first round.
+ ///
+ internal static TimeSpan StartupOffset(TimeSpan idleInterval, double sample)
+ => TimeSpan.FromTicks((long)(idleInterval.Ticks * sample));
+
+ ///
+ internal static TimeSpan StartupOffset(TimeSpan idleInterval)
+ => StartupOffset(idleInterval, Random.Shared.NextDouble());
+
+ private static TimeSpan Double(TimeSpan value) => TimeSpan.FromTicks(value.Ticks * 2);
+
+ private static TimeSpan MinOf(TimeSpan a, TimeSpan b) => a < b ? a : b;
+
+ private static TimeSpan MaxOf(TimeSpan a, TimeSpan b) => a > b ? a : b;
+}
diff --git a/framework/test/BBT.Aether.Infrastructure.Tests/BBT/Aether/Events/Processing/OutboxBackgroundServiceTests.cs b/framework/test/BBT.Aether.Infrastructure.Tests/BBT/Aether/Events/Processing/OutboxBackgroundServiceTests.cs
index 08f42c7..9165011 100644
--- a/framework/test/BBT.Aether.Infrastructure.Tests/BBT/Aether/Events/Processing/OutboxBackgroundServiceTests.cs
+++ b/framework/test/BBT.Aether.Infrastructure.Tests/BBT/Aether/Events/Processing/OutboxBackgroundServiceTests.cs
@@ -1,51 +1,120 @@
using System;
using BBT.Aether.Events;
+using BBT.Aether.Polling;
using Shouldly;
using Xunit;
namespace BBT.Aether.Events.Processing;
+///
+/// Pins the adaptive poll pacing that the outbox and inbox loops share. These exercise
+/// itself — the production code — rather than a copy of its arithmetic, so
+/// a change to the pacing rules cannot pass unnoticed.
+///
public sealed class AdaptivePollingTests
{
- private static TimeSpan NextDelay(TimeSpan current, int processed, AetherOutboxOptions opts)
+ private static AetherOutboxOptions Options() => new()
{
- if (processed > 0) return opts.BusyPollingInterval;
- var next = TimeSpan.FromMilliseconds(current.TotalMilliseconds * 2);
- return next > opts.MaxPollingInterval ? opts.MaxPollingInterval : next;
- }
+ BusyPollingInterval = TimeSpan.FromMilliseconds(100),
+ IdlePollingInterval = TimeSpan.FromSeconds(5),
+ MaxPollingInterval = TimeSpan.FromSeconds(60),
+ };
[Fact]
public void Busy_returns_busy_interval()
{
- var opts = new AetherOutboxOptions
- {
- BusyPollingInterval = TimeSpan.FromMilliseconds(100),
- IdlePollingInterval = TimeSpan.FromSeconds(5),
- MaxPollingInterval = TimeSpan.FromSeconds(60),
- };
- NextDelay(opts.IdlePollingInterval, processed: 10, opts)
- .ShouldBe(opts.BusyPollingInterval);
+ var opts = Options();
+ PollingDelay.OnProcessed(opts.BusyPollingInterval).ShouldBe(opts.BusyPollingInterval);
}
[Fact]
- public void Idle_doubles_delay_each_round()
+ public void Idle_doubles_delay_each_round_and_caps()
{
- var opts = new AetherOutboxOptions
- {
- BusyPollingInterval = TimeSpan.FromMilliseconds(100),
- IdlePollingInterval = TimeSpan.FromSeconds(5),
- MaxPollingInterval = TimeSpan.FromSeconds(60),
- };
- var d1 = NextDelay(opts.IdlePollingInterval, processed: 0, opts); // 10s
- var d2 = NextDelay(d1, processed: 0, opts); // 20s
- var d3 = NextDelay(d2, processed: 0, opts); // 40s
- var d4 = NextDelay(d3, processed: 0, opts); // 60s (capped)
- var d5 = NextDelay(d4, processed: 0, opts); // 60s (stays capped)
+ var opts = Options();
+ var d1 = PollingDelay.OnEmpty(opts.IdlePollingInterval, opts.MaxPollingInterval); // 10s
+ var d2 = PollingDelay.OnEmpty(d1, opts.MaxPollingInterval); // 20s
+ var d3 = PollingDelay.OnEmpty(d2, opts.MaxPollingInterval); // 40s
+ var d4 = PollingDelay.OnEmpty(d3, opts.MaxPollingInterval); // 60s capped
+ var d5 = PollingDelay.OnEmpty(d4, opts.MaxPollingInterval); // stays capped
d1.ShouldBe(TimeSpan.FromSeconds(10));
d2.ShouldBe(TimeSpan.FromSeconds(20));
d3.ShouldBe(TimeSpan.FromSeconds(40));
- d4.ShouldBe(TimeSpan.FromSeconds(60));
- d5.ShouldBe(TimeSpan.FromSeconds(60));
+ d4.ShouldBe(opts.MaxPollingInterval);
+ d5.ShouldBe(opts.MaxPollingInterval);
+ }
+
+ [Fact]
+ public void Error_backs_off_one_step_instead_of_jumping_to_the_cap()
+ {
+ // The old behaviour set the delay to MaxPollingInterval on any exception, so one transient
+ // fault stalled every replica for a full minute. Escalation keeps a blip cheap.
+ var opts = Options();
+
+ var first = PollingDelay.OnError(opts.BusyPollingInterval, opts.IdlePollingInterval, opts.MaxPollingInterval);
+
+ first.ShouldBe(opts.IdlePollingInterval);
+ first.ShouldBeLessThan(opts.MaxPollingInterval);
+ }
+
+ [Fact]
+ public void Error_never_retries_at_the_busy_cadence()
+ {
+ var opts = Options();
+
+ // Straight after a busy round the delay is 100 ms; doubling alone would retry a hard failure
+ // 5 times a second, so the idle interval is the floor.
+ PollingDelay.OnError(TimeSpan.FromMilliseconds(100), opts.IdlePollingInterval, opts.MaxPollingInterval)
+ .ShouldBeGreaterThanOrEqualTo(opts.IdlePollingInterval);
+ }
+
+ [Fact]
+ public void Repeated_errors_still_escalate_to_the_cap()
+ {
+ var opts = Options();
+ var d = opts.BusyPollingInterval;
+ for (var i = 0; i < 10; i++)
+ d = PollingDelay.OnError(d, opts.IdlePollingInterval, opts.MaxPollingInterval);
+
+ d.ShouldBe(opts.MaxPollingInterval);
+ }
+
+ [Theory]
+ [InlineData(0.0, 0.75)]
+ [InlineData(0.5, 1.00)]
+ [InlineData(1.0, 1.25)]
+ public void Jitter_spans_the_configured_fraction_either_side(double sample, double expectedScale)
+ {
+ var nominal = TimeSpan.FromSeconds(60);
+
+ var jittered = PollingDelay.Jitter(nominal, sample);
+
+ jittered.TotalSeconds.ShouldBe(60 * expectedScale, tolerance: 0.001);
+ }
+
+ [Fact]
+ public void Jitter_never_returns_a_non_positive_delay()
+ {
+ PollingDelay.Jitter(TimeSpan.Zero, 0.0).ShouldBeGreaterThan(TimeSpan.Zero);
+ PollingDelay.Jitter(TimeSpan.FromTicks(1), 0.0).ShouldBeGreaterThan(TimeSpan.Zero);
+ }
+
+ [Fact]
+ public void Jitter_keeps_replicas_from_sharing_a_phase()
+ {
+ // Two replicas holding the same nominal delay must not wake together.
+ var nominal = TimeSpan.FromSeconds(60);
+
+ PollingDelay.Jitter(nominal, 0.1).ShouldNotBe(PollingDelay.Jitter(nominal, 0.9));
+ }
+
+ [Fact]
+ public void Startup_offset_stays_within_the_idle_interval()
+ {
+ var opts = Options();
+
+ PollingDelay.StartupOffset(opts.IdlePollingInterval, 0.0).ShouldBe(TimeSpan.Zero);
+ PollingDelay.StartupOffset(opts.IdlePollingInterval, 0.999)
+ .ShouldBeLessThan(opts.IdlePollingInterval);
}
}