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))); + } +}