From c7912fd2c81cb78086c479345f756b8f8f9ae38e Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 21 Jun 2026 01:28:40 +0200 Subject: [PATCH 1/2] Fix scheduling startup backlog catch-up Rebuild local schedules in bounded pages and stagger past-due specific-instant catch-up so orphaned scheduling bookmarks do not flood dispatch during startup. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../Modules/Runtime/BookmarkStore.cs | 10 ++++ .../Modules/Runtime/TriggerStore.cs | 5 +- .../Features/SchedulingFeature.cs | 3 + .../Options/SchedulingOptions.cs | 27 +++++++++ .../ScheduledSpecificInstantTask.cs | 36 +++++++++-- .../Services/PastDueScheduleStaggerer.cs | 34 +++++++++++ .../ShellFeatures/SchedulingFeature.cs | 3 + .../CreateSchedulesStartupTask.cs | 56 ++++++++++++++++-- .../Contracts/IBookmarkStore.cs | 21 +++++++ .../Stores/CachingTriggerStore.cs | 4 +- .../Stores/MemoryBookmarkStore.cs | 10 ++++ .../HttpWorkflowsMiddlewareTests.cs | 8 +++ .../ScheduledSpecificInstantTaskTests.cs | 59 +++++++++++++++---- .../CreateSchedulesStartupTaskTests.cs | 48 +++++++++++++-- 14 files changed, 290 insertions(+), 34 deletions(-) create mode 100644 src/modules/Elsa.Scheduling/Options/SchedulingOptions.cs create mode 100644 src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs diff --git a/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/BookmarkStore.cs b/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/BookmarkStore.cs index 8d421a6b8..cb50bfca9 100644 --- a/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/BookmarkStore.cs +++ b/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/BookmarkStore.cs @@ -1,3 +1,5 @@ +using Elsa.Common.Models; +using Elsa.Extensions; using Elsa.Workflows; using Elsa.Workflows.Runtime; using Elsa.Workflows.Runtime.Entities; @@ -36,6 +38,14 @@ public class EFCoreBookmarkStore(Store sto return await store.QueryAsync(filter.Apply, OnLoadAsync, filter.TenantAgnostic, cancellationToken); } + /// + public async ValueTask> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) + { + var count = await store.CountAsync(filter.Apply, filter.TenantAgnostic, cancellationToken); + var results = (await store.QueryAsync(query => filter.Apply(query).OrderBy(x => x.Id).Paginate(pageArgs), OnLoadAsync, filter.TenantAgnostic, cancellationToken)).ToList(); + return Page.Of(results, count); + } + /// public async ValueTask DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/TriggerStore.cs b/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/TriggerStore.cs index b160b94e9..1e0f8cd3a 100644 --- a/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/TriggerStore.cs +++ b/src/modules/Elsa.Persistence.EFCore/Modules/Runtime/TriggerStore.cs @@ -8,7 +8,6 @@ using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.OrderDefinitions; using JetBrains.Annotations; -using Open.Linq.AsyncExtensions; namespace Elsa.Persistence.EFCore.Modules.Runtime; @@ -50,8 +49,8 @@ public class EFCoreTriggerStore( public async ValueTask> FindManyAsync(TriggerFilter filter, PageArgs pageArgs, StoredTriggerOrder order, CancellationToken cancellationToken = default) { - var count = await store.QueryAsync(filter.Apply, OnLoadAsync, cancellationToken).LongCount(); - var results = await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(order).Paginate(pageArgs).OrderBy(order), OnLoadAsync, cancellationToken).ToList(); + var count = await store.CountAsync(filter.Apply, filter.TenantAgnostic, cancellationToken); + var results = (await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(order).Paginate(pageArgs).OrderBy(order), OnLoadAsync, filter.TenantAgnostic, cancellationToken)).ToList(); return new(results, count); } diff --git a/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs index 20efd0a61..9f17cffed 100644 --- a/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs +++ b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs @@ -6,6 +6,7 @@ using Elsa.Features.Attributes; using Elsa.Features.Services; using Elsa.Scheduling.Bookmarks; using Elsa.Scheduling.Handlers; +using Elsa.Scheduling.Options; using Elsa.Scheduling.Services; using Elsa.Scheduling.StartupTasks; using Elsa.Scheduling.TriggerPayloadValidators; @@ -44,6 +45,7 @@ public class SchedulingFeature : FeatureBase .AddSingleton() .AddSingleton(sp => sp.GetRequiredService()) .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton(CronParser) .AddScoped() @@ -59,6 +61,7 @@ public class SchedulingFeature : FeatureBase // Graceful shutdown: register scheduled-trigger ingress for diagnostic visibility (FR-006). .AddSingleton(); + Services.Configure(_ => { }); Services.Configure(options => { options.RegisterTypeAlias(typeof(CronBookmarkPayload), nameof(CronBookmarkPayload)); diff --git a/src/modules/Elsa.Scheduling/Options/SchedulingOptions.cs b/src/modules/Elsa.Scheduling/Options/SchedulingOptions.cs new file mode 100644 index 000000000..efd78fc84 --- /dev/null +++ b/src/modules/Elsa.Scheduling/Options/SchedulingOptions.cs @@ -0,0 +1,27 @@ +namespace Elsa.Scheduling.Options; + +/// +/// Configures local scheduling behavior. +/// +public class SchedulingOptions +{ + /// + /// The number of stored triggers and bookmarks to load per batch when rebuilding local schedules on startup. + /// + public int StartupSchedulePageSize { get; set; } = 1000; + + /// + /// The minimum delay used when a specific-instant schedule is already due. + /// + public TimeSpan MinimumPastDueScheduleDelay { get; set; } = TimeSpan.FromMilliseconds(1); + + /// + /// The spacing between past-due schedules during catch-up. + /// + public TimeSpan PastDueScheduleStaggerInterval { get; set; } = TimeSpan.FromMilliseconds(50); + + /// + /// The bounded window over which past-due schedules are distributed. + /// + public TimeSpan PastDueScheduleStaggerWindow { get; set; } = TimeSpan.FromMinutes(5); +} diff --git a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs index 51df7f954..d081290b4 100644 --- a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs +++ b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs @@ -1,9 +1,12 @@ using Elsa.Common; using Elsa.Mediator.Contracts; using Elsa.Scheduling.Commands; +using Elsa.Scheduling.Options; +using Elsa.Scheduling.Services; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Timer = System.Timers.Timer; +using OptionsFactory = Microsoft.Extensions.Options.Options; namespace Elsa.Scheduling.ScheduledTasks; @@ -12,10 +15,12 @@ namespace Elsa.Scheduling.ScheduledTasks; /// public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable { + private static readonly PastDueScheduleStaggerer DefaultPastDueScheduleStaggerer = new(OptionsFactory.Create(new SchedulingOptions())); private readonly ITask _task; private readonly ISystemClock _systemClock; private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; + private readonly PastDueScheduleStaggerer _pastDueScheduleStaggerer; private readonly DateTimeOffset _startAt; private readonly CancellationTokenSource _cancellationTokenSource; private readonly SemaphoreSlim _executionSemaphore = new(1, 1); @@ -27,12 +32,33 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable /// /// Initializes a new instance of . /// - public ScheduledSpecificInstantTask(ITask task, DateTimeOffset startAt, ISystemClock systemClock, IServiceScopeFactory scopeFactory, ILogger logger) + public ScheduledSpecificInstantTask( + ITask task, + DateTimeOffset startAt, + ISystemClock systemClock, + IServiceScopeFactory scopeFactory, + ILogger logger) + : this(task, startAt, systemClock, scopeFactory, logger, DefaultPastDueScheduleStaggerer) + { + } + + /// + /// Initializes a new instance of . + /// + [ActivatorUtilitiesConstructor] + public ScheduledSpecificInstantTask( + ITask task, + DateTimeOffset startAt, + ISystemClock systemClock, + IServiceScopeFactory scopeFactory, + ILogger logger, + PastDueScheduleStaggerer pastDueScheduleStaggerer) { _task = task; _systemClock = systemClock; _scopeFactory = scopeFactory; _logger = logger; + _pastDueScheduleStaggerer = pastDueScheduleStaggerer; _startAt = startAt; _cancellationTokenSource = new(); @@ -57,16 +83,14 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable { var now = _systemClock.UtcNow; var delay = _startAt - now; + var adjustedDelay = _pastDueScheduleStaggerer.GetDelay(delay); - // Handle edge cases where delay is zero or negative (e.g., due to clock drift, fast execution, or time alignment) - // Instead of silently returning, use a minimum delay to ensure the timer fires and workflow continues scheduling if (delay <= TimeSpan.Zero) { - _logger.LogWarning("Calculated delay is {Delay} which is not positive. Using minimum delay of 1ms to ensure timer fires", delay); - delay = TimeSpan.FromMilliseconds(1); + _logger.LogDebug("Calculated delay is {Delay} which is not positive. Using catch-up delay of {CatchUpDelay}", delay, adjustedDelay); } - _timer = new(delay.TotalMilliseconds) + _timer = new(adjustedDelay.TotalMilliseconds) { Enabled = true }; diff --git a/src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs b/src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs new file mode 100644 index 000000000..a5d72d53b --- /dev/null +++ b/src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs @@ -0,0 +1,34 @@ +using Elsa.Scheduling.Options; +using Microsoft.Extensions.Options; + +namespace Elsa.Scheduling.Services; + +/// +/// Distributes already-due schedules over a bounded window to avoid dispatch storms during startup catch-up. +/// +public class PastDueScheduleStaggerer(IOptions options) +{ + private long _sequence; + + public TimeSpan GetDelay(TimeSpan calculatedDelay) + { + if (calculatedDelay > TimeSpan.Zero) + return calculatedDelay; + + var currentOptions = options.Value; + var minimumDelay = GetPositiveOrDefault(currentOptions.MinimumPastDueScheduleDelay, TimeSpan.FromMilliseconds(1)); + var staggerInterval = currentOptions.PastDueScheduleStaggerInterval; + var staggerWindow = currentOptions.PastDueScheduleStaggerWindow; + + if (staggerInterval <= TimeSpan.Zero || staggerWindow <= TimeSpan.Zero) + return minimumDelay; + + var slotCount = Math.Max(1, staggerWindow.Ticks / staggerInterval.Ticks); + var sequence = Interlocked.Increment(ref _sequence) - 1; + var slot = (sequence & long.MaxValue) % slotCount; + + return minimumDelay + TimeSpan.FromTicks(staggerInterval.Ticks * slot); + } + + private static TimeSpan GetPositiveOrDefault(TimeSpan value, TimeSpan defaultValue) => value > TimeSpan.Zero ? value : defaultValue; +} diff --git a/src/modules/Elsa.Scheduling/ShellFeatures/SchedulingFeature.cs b/src/modules/Elsa.Scheduling/ShellFeatures/SchedulingFeature.cs index 5dd3dfb59..6f22ee2ba 100644 --- a/src/modules/Elsa.Scheduling/ShellFeatures/SchedulingFeature.cs +++ b/src/modules/Elsa.Scheduling/ShellFeatures/SchedulingFeature.cs @@ -5,6 +5,7 @@ using Elsa.Extensions; using Elsa.Workflows.Options; using Elsa.Scheduling.Bookmarks; using Elsa.Scheduling.Handlers; +using Elsa.Scheduling.Options; using Elsa.Scheduling.Services; using Elsa.Scheduling.StartupTasks; using Elsa.Scheduling.TriggerPayloadValidators; @@ -44,6 +45,7 @@ public class SchedulingFeature : IShellFeature .AddSingleton() .AddSingleton(sp => sp.GetRequiredService()) .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton(CronParser) .AddScoped() @@ -55,6 +57,7 @@ public class SchedulingFeature : IShellFeature .AddTriggerPayloadValidator() .AddActivitiesFrom(); + services.Configure(_ => { }); services.Configure(options => { options.AddTypeAlias(); diff --git a/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs b/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs index 4ab8971d6..7b4c1caeb 100644 --- a/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs +++ b/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs @@ -1,9 +1,12 @@ using Elsa.Common; using Elsa.Common.Multitenancy; +using Elsa.Common.Models; +using Elsa.Scheduling.Options; using Elsa.Workflows.Runtime; using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.Tasks; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; namespace Elsa.Scheduling.StartupTasks; @@ -11,7 +14,7 @@ namespace Elsa.Scheduling.StartupTasks; /// Enqueues schedule creation when using the default scheduler, which doesn't have its own persistence layer like Quartz or Hangfire. /// [TaskDependency(typeof(PopulateRegistriesStartupTask))] -public class CreateSchedulesStartupTask(IServiceProvider serviceProvider) : IStartupTask +public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptions options) : IStartupTask { public async Task ExecuteAsync(CancellationToken cancellationToken) { @@ -23,12 +26,13 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider) : ISta await CreateSchedulesAsync(serviceProvider, cancellationToken); } - private static async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) + private async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) { var triggerStore = serviceProvider.GetRequiredService(); var bookmarkStore = serviceProvider.GetRequiredService(); var triggerScheduler = serviceProvider.GetRequiredService(); var bookmarkScheduler = serviceProvider.GetRequiredService(); + var pageSize = Math.Max(1, options.Value.StartupSchedulePageSize); var stimulusNames = new[] { SchedulingStimulusNames.Cron, SchedulingStimulusNames.Timer, SchedulingStimulusNames.StartAt, SchedulingStimulusNames.Delay, @@ -41,10 +45,50 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider) : ISta { Names = stimulusNames }; - var triggers = (await triggerStore.FindManyAsync(triggerFilter, cancellationToken)).ToList(); - var bookmarks = (await bookmarkStore.FindManyAsync(bookmarkFilter, cancellationToken)).ToList(); - await triggerScheduler.ScheduleAsync(triggers, cancellationToken); - await bookmarkScheduler.ScheduleAsync(bookmarks, cancellationToken); + await ScheduleTriggersAsync(triggerStore, triggerScheduler, triggerFilter, pageSize, cancellationToken); + await ScheduleBookmarksAsync(bookmarkStore, bookmarkScheduler, bookmarkFilter, pageSize, cancellationToken); + } + + private static async Task ScheduleTriggersAsync(ITriggerStore triggerStore, ITriggerScheduler triggerScheduler, TriggerFilter triggerFilter, int pageSize, CancellationToken cancellationToken) + { + var pageArgs = PageArgs.FromRange(0, pageSize); + + while (true) + { + var page = await triggerStore.FindManyAsync(triggerFilter, pageArgs, cancellationToken); + + if (page.Items.Count == 0) + break; + + await triggerScheduler.ScheduleAsync(page.Items, cancellationToken); + + var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count; + if (nextOffset >= page.TotalCount) + break; + + pageArgs = pageArgs.Next(); + } + } + + private static async Task ScheduleBookmarksAsync(IBookmarkStore bookmarkStore, IBookmarkScheduler bookmarkScheduler, BookmarkFilter bookmarkFilter, int pageSize, CancellationToken cancellationToken) + { + var pageArgs = PageArgs.FromRange(0, pageSize); + + while (true) + { + var page = await bookmarkStore.FindManyAsync(bookmarkFilter, pageArgs, cancellationToken); + + if (page.Items.Count == 0) + break; + + await bookmarkScheduler.ScheduleAsync(page.Items, cancellationToken); + + var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count; + if (nextOffset >= page.TotalCount) + break; + + pageArgs = pageArgs.Next(); + } } } diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs index 679090cdd..a6c07c173 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs @@ -1,3 +1,4 @@ +using Elsa.Common.Models; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; @@ -34,6 +35,26 @@ public interface IBookmarkStore /// ValueTask> FindManyAsync(BookmarkFilter filter, CancellationToken cancellationToken = default); + /// + /// Returns a page of bookmarks matching the specified filter. + /// + /// + /// The default implementation materializes all matching bookmarks and pages them in memory. Stores backed by external persistence should override this method. + /// + async ValueTask> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) + { + var records = (await FindManyAsync(filter, cancellationToken)).OrderBy(x => x.Id).ToList(); + IEnumerable page = records; + + if (pageArgs.Offset.HasValue) + page = page.Skip(pageArgs.Offset.Value); + + if (pageArgs.Limit.HasValue) + page = page.Take(pageArgs.Limit.Value); + + return Page.Of(page.ToList(), records.Count); + } + /// /// Deletes a set of bookmarks matching the specified filter. /// diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs index e84845dc9..50b7bdade 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs @@ -48,13 +48,13 @@ public class CachingTriggerStore(ITriggerStore decoratedStore, ICacheManager cac public async ValueTask> FindManyAsync(TriggerFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) { - var cacheKey = hasher.Hash(filter); + var cacheKey = hasher.Hash(filter, pageArgs); return (await GetOrCreateAsync(cacheKey, async () => await decoratedStore.FindManyAsync(filter, pageArgs, cancellationToken)))!; } public async ValueTask> FindManyAsync(TriggerFilter filter, PageArgs pageArgs, StoredTriggerOrder order, CancellationToken cancellationToken = default) { - var cacheKey = hasher.Hash(filter); + var cacheKey = hasher.Hash(filter, pageArgs, order); return (await GetOrCreateAsync(cacheKey, async () => await decoratedStore.FindManyAsync(filter, pageArgs, order, cancellationToken)))!; } diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkStore.cs index 6cb8b37d8..b2115e28e 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkStore.cs @@ -1,4 +1,6 @@ +using Elsa.Common.Models; using Elsa.Common.Services; +using Elsa.Extensions; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; using JetBrains.Annotations; @@ -37,6 +39,14 @@ public class MemoryBookmarkStore(MemoryStore store) : IBookmarkS return new(entities); } + /// + public ValueTask> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) + { + var count = store.Query(query => Filter(query, filter)).LongCount(); + var result = store.Query(query => Filter(query, filter).OrderBy(x => x.Id).Paginate(pageArgs)).ToList(); + return new(Page.Of(result, count)); + } + /// public async ValueTask DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default) { diff --git a/test/unit/Elsa.Http.UnitTests/Middleware/HttpWorkflowsMiddlewareTests.cs b/test/unit/Elsa.Http.UnitTests/Middleware/HttpWorkflowsMiddlewareTests.cs index 529b3c226..4d2f9616d 100644 --- a/test/unit/Elsa.Http.UnitTests/Middleware/HttpWorkflowsMiddlewareTests.cs +++ b/test/unit/Elsa.Http.UnitTests/Middleware/HttpWorkflowsMiddlewareTests.cs @@ -1,4 +1,5 @@ using System.Collections; +using Elsa.Common.Models; using Elsa.Http.Bookmarks; using Elsa.Http.Middleware; using Elsa.Http.Options; @@ -104,6 +105,13 @@ public class HttpWorkflowsMiddlewareTests return new([]); } + public ValueTask> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) + { + LastFilter = filter; + var results = Filter(filter).ToList(); + return new(Page.Of(results, results.Count)); + } + public ValueTask DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default) { var bookmarksToDelete = Filter(filter).ToList(); diff --git a/test/unit/Elsa.Scheduling.UnitTests/ScheduledTasks/ScheduledSpecificInstantTaskTests.cs b/test/unit/Elsa.Scheduling.UnitTests/ScheduledTasks/ScheduledSpecificInstantTaskTests.cs index 8321d796b..c4052e931 100644 --- a/test/unit/Elsa.Scheduling.UnitTests/ScheduledTasks/ScheduledSpecificInstantTaskTests.cs +++ b/test/unit/Elsa.Scheduling.UnitTests/ScheduledTasks/ScheduledSpecificInstantTaskTests.cs @@ -1,9 +1,12 @@ using Elsa.Common; using Elsa.Mediator.Contracts; +using Elsa.Scheduling.Options; using Elsa.Scheduling.ScheduledTasks; +using Elsa.Scheduling.Services; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using NSubstitute; +using OptionsFactory = Microsoft.Extensions.Options.Options; namespace Elsa.Scheduling.UnitTests.ScheduledTasks; @@ -17,6 +20,12 @@ public class ScheduledSpecificInstantTaskTests : IDisposable private readonly ServiceProvider _serviceProvider; private readonly ISystemClock _systemClock; private readonly ILogger _logger; + private readonly PastDueScheduleStaggerer _pastDueScheduleStaggerer = new(OptionsFactory.Create(new SchedulingOptions + { + MinimumPastDueScheduleDelay = TimeSpan.FromMilliseconds(10), + PastDueScheduleStaggerInterval = TimeSpan.FromMilliseconds(25), + PastDueScheduleStaggerWindow = TimeSpan.FromMilliseconds(100) + })); private readonly List _tasksToDispose = new(); public ScheduledSpecificInstantTaskTests() @@ -39,7 +48,8 @@ public class ScheduledSpecificInstantTaskTests : IDisposable startAt ?? DefaultNow.AddMinutes(5), systemClock ?? _systemClock, _serviceProvider.CreateScope().ServiceProvider.GetRequiredService(), - _logger + _logger, + _pastDueScheduleStaggerer ); _tasksToDispose.Add(scheduledTask); return scheduledTask; @@ -60,10 +70,10 @@ public class ScheduledSpecificInstantTaskTests : IDisposable Arg.Any>()); } - private void AssertWarningLogged(int expectedCount = 1) + private void AssertDebugLogged(int expectedCount = 1) { _logger.Received(expectedCount).Log( - LogLevel.Warning, + LogLevel.Debug, Arg.Any(), Arg.Any(), Arg.Any(), @@ -85,31 +95,58 @@ public class ScheduledSpecificInstantTaskTests : IDisposable } [Fact] - public void Schedule_WithZeroDelay_ShouldUseMinimumDelay() + public void Schedule_WithZeroDelay_ShouldUseCatchUpDelay() { // Arrange - simulate a case where startAt is exactly now SetupSystemClock(DefaultNow); var startAt = DefaultNow; // delay = 0 - // Act - Should adjust to 1ms minimum delay + // Act - Should adjust to a bounded catch-up delay CreateScheduledTask(startAt: startAt); - // Assert - Should not crash and should log warning - AssertWarningLogged(); + AssertDebugLogged(); } [Fact] - public void Schedule_WithNegativeDelay_ShouldUseMinimumDelay() + public void Schedule_WithNegativeDelay_ShouldUseCatchUpDelay() { // Arrange - simulate a case where startAt is in the past SetupSystemClock(DefaultNow); var startAt = DefaultNow.AddMinutes(-1); // Past time - // Act - Should adjust to 1ms minimum delay + // Act - Should adjust to a bounded catch-up delay CreateScheduledTask(startAt: startAt); - // Assert - Should log a warning - AssertWarningLogged(); + AssertDebugLogged(); + } + + [Fact] + public void Schedule_WithMultiplePastDueTasks_ShouldStaggerCatchUp() + { + SetupSystemClock(DefaultNow); + + CreateScheduledTask(startAt: DefaultNow.AddMinutes(-1)); + CreateScheduledTask(startAt: DefaultNow.AddMinutes(-1)); + CreateScheduledTask(startAt: DefaultNow.AddMinutes(-1)); + + _logger.Received(1).Log( + LogLevel.Debug, + Arg.Any(), + Arg.Is(x => x.ToString()!.Contains("00:00:00.0100000")), + Arg.Any(), + Arg.Any>()); + _logger.Received(1).Log( + LogLevel.Debug, + Arg.Any(), + Arg.Is(x => x.ToString()!.Contains("00:00:00.0350000")), + Arg.Any(), + Arg.Any>()); + _logger.Received(1).Log( + LogLevel.Debug, + Arg.Any(), + Arg.Is(x => x.ToString()!.Contains("00:00:00.0600000")), + Arg.Any(), + Arg.Any>()); } [Fact] diff --git a/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs b/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs index d6c16a3a6..cdeee7dbe 100644 --- a/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs +++ b/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs @@ -1,5 +1,7 @@ using Elsa.Common; using Elsa.Common.Multitenancy; +using Elsa.Common.Models; +using Elsa.Scheduling.Options; using Elsa.Scheduling.StartupTasks; using Elsa.Workflows.Runtime; using Elsa.Workflows.Runtime.Entities; @@ -7,22 +9,30 @@ using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.Tasks; using Microsoft.Extensions.DependencyInjection; using NSubstitute; +using OptionsFactory = Microsoft.Extensions.Options.Options; namespace Elsa.Scheduling.UnitTests.StartupTasks; public class CreateSchedulesStartupTaskTests { - private readonly StoredTrigger[] _triggers = [new() { WorkflowDefinitionId = "definition", WorkflowDefinitionVersionId = "version", ActivityId = "activity" }]; - private readonly StoredBookmark[] _bookmarks = [new() { Hash = "hash", WorkflowInstanceId = "instance" }]; + private readonly StoredTrigger[] _triggers = + [ + new() { Id = "trigger-1", WorkflowDefinitionId = "definition", WorkflowDefinitionVersionId = "version", ActivityId = "activity" } + ]; + + private readonly StoredBookmark[] _bookmarks = [new() { Id = "bookmark-1", Hash = "hash", WorkflowInstanceId = "instance" }]; private readonly ITriggerStore _triggerStore = Substitute.For(); private readonly IBookmarkStore _bookmarkStore = Substitute.For(); private readonly ITriggerScheduler _triggerScheduler = Substitute.For(); private readonly IBookmarkScheduler _bookmarkScheduler = Substitute.For(); + private readonly SchedulingOptions _options = new() { StartupSchedulePageSize = 1000 }; public CreateSchedulesStartupTaskTests() { - _triggerStore.FindManyAsync(Arg.Any(), Arg.Any()).Returns(_triggers); - _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any()).Returns(_bookmarks); + _triggerStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(new Page(_triggers, _triggers.Length)); + _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(new Page(_bookmarks, _bookmarks.Length)); } [Fact] @@ -36,7 +46,7 @@ public class CreateSchedulesStartupTaskTests [Fact] public async Task ExecuteAsync_WithoutTenantBackgroundQueue_SchedulesImmediately() { - var task = new CreateSchedulesStartupTask(CreateServiceProvider()); + var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); @@ -51,7 +61,7 @@ public class CreateSchedulesStartupTaskTests var workQueue = Substitute.For(); workQueue.EnqueueAsync(Arg.Do(x => workItem = x), Arg.Any()).Returns(ValueTask.CompletedTask); var serviceProvider = CreateServiceProvider(services => services.AddSingleton(workQueue)); - var task = new CreateSchedulesStartupTask(serviceProvider); + var task = new CreateSchedulesStartupTask(serviceProvider, OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); @@ -65,6 +75,32 @@ public class CreateSchedulesStartupTaskTests await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(_bookmarks)), Arg.Any()); } + [Fact] + public async Task ExecuteAsync_SchedulesInConfiguredPages() + { + var firstTriggerPage = new[] { _triggers[0] }; + var secondTriggerPage = new[] { new StoredTrigger { Id = "trigger-2", WorkflowDefinitionId = "definition", WorkflowDefinitionVersionId = "version", ActivityId = "activity" } }; + var firstBookmarkPage = new[] { _bookmarks[0] }; + var secondBookmarkPage = new[] { new StoredBookmark { Id = "bookmark-2", Hash = "hash", WorkflowInstanceId = "instance" } }; + _options.StartupSchedulePageSize = 1; + _triggerStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 0 && x.Limit == 1), Arg.Any()) + .Returns(new Page(firstTriggerPage, 2)); + _triggerStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 1 && x.Limit == 1), Arg.Any()) + .Returns(new Page(secondTriggerPage, 2)); + _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 0 && x.Limit == 1), Arg.Any()) + .Returns(new Page(firstBookmarkPage, 2)); + _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 1 && x.Limit == 1), Arg.Any()) + .Returns(new Page(secondBookmarkPage, 2)); + var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); + + await task.ExecuteAsync(CancellationToken.None); + + await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(firstTriggerPage)), Arg.Any()); + await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(secondTriggerPage)), Arg.Any()); + await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(firstBookmarkPage)), Arg.Any()); + await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(secondBookmarkPage)), Arg.Any()); + } + private ServiceProvider CreateServiceProvider(Action? configureServices = null) { var services = new ServiceCollection(); From 7449d2e703d99cc14b8fe741a866953dc81f1264 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 21 Jun 2026 01:37:57 +0200 Subject: [PATCH 2/2] address greptile review feedback (greploop iteration 1) Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../Services/DefaultTriggerScheduler.cs | 10 ++--- .../Services/PastDueScheduleStaggerer.cs | 7 +++- .../Contracts/IBookmarkStore.cs | 16 +------ .../Services/DefaultTriggerSchedulerTests.cs | 42 +++++++++++++++++++ .../Services/PastDueScheduleStaggererTests.cs | 23 ++++++++++ 5 files changed, 76 insertions(+), 22 deletions(-) create mode 100644 test/unit/Elsa.Scheduling.UnitTests/Services/DefaultTriggerSchedulerTests.cs create mode 100644 test/unit/Elsa.Scheduling.UnitTests/Services/PastDueScheduleStaggererTests.cs diff --git a/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs b/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs index 71ce81450..434447b35 100644 --- a/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs +++ b/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs @@ -41,14 +41,10 @@ public class DefaultTriggerScheduler(IWorkflowScheduler workflowScheduler, ISyst foreach (var trigger in startAtTriggers) { var executeAt = trigger.GetPayload().ExecuteAt; - - // If the trigger is in the past, log info and skip scheduling. + if (executeAt < now) - { - logger.LogInformation("StartAt trigger is in the past. TriggerId: {TriggerId}. ExecuteAt: {ExecuteAt}. Skipping scheduling", trigger.Id, executeAt); - continue; - } - + logger.LogInformation("StartAt trigger is in the past. TriggerId: {TriggerId}. ExecuteAt: {ExecuteAt}. Scheduling catch-up", trigger.Id, executeAt); + var input = new { ExecuteAt = executeAt }.ToDictionary(); var request = new ScheduleNewWorkflowInstanceRequest { diff --git a/src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs b/src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs index a5d72d53b..6242c6165 100644 --- a/src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs +++ b/src/modules/Elsa.Scheduling/Services/PastDueScheduleStaggerer.cs @@ -23,7 +23,12 @@ public class PastDueScheduleStaggerer(IOptions options) if (staggerInterval <= TimeSpan.Zero || staggerWindow <= TimeSpan.Zero) return minimumDelay; - var slotCount = Math.Max(1, staggerWindow.Ticks / staggerInterval.Ticks); + var availableWindow = staggerWindow - minimumDelay; + + if (availableWindow <= TimeSpan.Zero) + return minimumDelay; + + var slotCount = Math.Max(1, availableWindow.Ticks / staggerInterval.Ticks + 1); var sequence = Interlocked.Increment(ref _sequence) - 1; var slot = (sequence & long.MaxValue) % slotCount; diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs index a6c07c173..013eabc72 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkStore.cs @@ -39,21 +39,9 @@ public interface IBookmarkStore /// Returns a page of bookmarks matching the specified filter. /// /// - /// The default implementation materializes all matching bookmarks and pages them in memory. Stores backed by external persistence should override this method. + /// Startup backlog catch-up depends on store-backed paging. Implementations should page at the persistence layer instead of materializing all matches in memory. /// - async ValueTask> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) - { - var records = (await FindManyAsync(filter, cancellationToken)).OrderBy(x => x.Id).ToList(); - IEnumerable page = records; - - if (pageArgs.Offset.HasValue) - page = page.Skip(pageArgs.Offset.Value); - - if (pageArgs.Limit.HasValue) - page = page.Take(pageArgs.Limit.Value); - - return Page.Of(page.ToList(), records.Count); - } + ValueTask> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default); /// /// Deletes a set of bookmarks matching the specified filter. diff --git a/test/unit/Elsa.Scheduling.UnitTests/Services/DefaultTriggerSchedulerTests.cs b/test/unit/Elsa.Scheduling.UnitTests/Services/DefaultTriggerSchedulerTests.cs new file mode 100644 index 000000000..b8376752c --- /dev/null +++ b/test/unit/Elsa.Scheduling.UnitTests/Services/DefaultTriggerSchedulerTests.cs @@ -0,0 +1,42 @@ +using Elsa.Common; +using Elsa.Scheduling.Activities; +using Elsa.Scheduling.Bookmarks; +using Elsa.Scheduling.Services; +using Elsa.Workflows.Helpers; +using Elsa.Workflows.Runtime.Entities; +using Microsoft.Extensions.Logging; +using NSubstitute; + +namespace Elsa.Scheduling.UnitTests.Services; + +public class DefaultTriggerSchedulerTests +{ + [Fact] + public async Task ScheduleAsync_SchedulesPastDueStartAtTriggerForCatchUp() + { + var workflowScheduler = Substitute.For(); + var systemClock = Substitute.For(); + var logger = Substitute.For>(); + var scheduler = new DefaultTriggerScheduler(workflowScheduler, systemClock, logger); + var now = new DateTimeOffset(2025, 11, 06, 22, 50, 00, TimeSpan.Zero); + var executeAt = now.AddMinutes(-5); + ScheduleNewWorkflowInstanceRequest? scheduledRequest = null; + var trigger = new StoredTrigger + { + Id = "trigger-1", + Name = ActivityTypeNameHelper.GenerateTypeName(), + WorkflowDefinitionVersionId = "workflow-version", + ActivityId = "activity-1", + Payload = new StartAtPayload(executeAt) + }; + systemClock.UtcNow.Returns(now); + workflowScheduler.ScheduleAtAsync(trigger.Id, Arg.Do(x => scheduledRequest = x), executeAt, Arg.Any()).Returns(ValueTask.CompletedTask); + + await scheduler.ScheduleAsync([trigger], CancellationToken.None); + + await workflowScheduler.Received(1).ScheduleAtAsync(trigger.Id, Arg.Any(), executeAt, Arg.Any()); + Assert.NotNull(scheduledRequest); + Assert.Equal(trigger.ActivityId, scheduledRequest.TriggerActivityId); + Assert.Equal(trigger.WorkflowDefinitionVersionId, scheduledRequest.WorkflowDefinitionHandle.DefinitionVersionId); + } +} diff --git a/test/unit/Elsa.Scheduling.UnitTests/Services/PastDueScheduleStaggererTests.cs b/test/unit/Elsa.Scheduling.UnitTests/Services/PastDueScheduleStaggererTests.cs new file mode 100644 index 000000000..263be2565 --- /dev/null +++ b/test/unit/Elsa.Scheduling.UnitTests/Services/PastDueScheduleStaggererTests.cs @@ -0,0 +1,23 @@ +using Elsa.Scheduling.Options; +using Elsa.Scheduling.Services; +using OptionsFactory = Microsoft.Extensions.Options.Options; + +namespace Elsa.Scheduling.UnitTests.Services; + +public class PastDueScheduleStaggererTests +{ + [Fact] + public void GetDelay_DoesNotExceedConfiguredWindow() + { + var staggerer = new PastDueScheduleStaggerer(OptionsFactory.Create(new SchedulingOptions + { + MinimumPastDueScheduleDelay = TimeSpan.FromSeconds(1), + PastDueScheduleStaggerInterval = TimeSpan.FromMilliseconds(900), + PastDueScheduleStaggerWindow = TimeSpan.FromSeconds(5) + })); + + var delays = Enumerable.Range(0, 16).Select(_ => staggerer.GetDelay(TimeSpan.Zero)).ToList(); + + Assert.All(delays, delay => Assert.InRange(delay, TimeSpan.FromSeconds(1), TimeSpan.FromSeconds(5))); + } +}