From 411151d748e178d6ae413dcbf2f1d2cb37ec5e2b Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 29 May 2024 10:34:44 +0200 Subject: [PATCH] Refactor `DefaultTriggerScheduler` and handle past `StartAt` triggers (#5463) The `DefaultTriggerScheduler` received a significant refactoring, with the injection of the `ISystemClock` service and a modification of several method calls. Additionally, a check is included to avoid scheduling `StartAt` triggers if their execution time is in the past. For these triggers, an information message is logged and scheduling is skipped. --- .../Elsa.Scheduling/Activities/StartAt.cs | 31 ++++++++-------- .../Services/DefaultTriggerScheduler.cs | 35 +++++++++---------- 2 files changed, 34 insertions(+), 32 deletions(-) diff --git a/src/modules/Elsa.Scheduling/Activities/StartAt.cs b/src/modules/Elsa.Scheduling/Activities/StartAt.cs index a773ff8b8..7b679acb7 100644 --- a/src/modules/Elsa.Scheduling/Activities/StartAt.cs +++ b/src/modules/Elsa.Scheduling/Activities/StartAt.cs @@ -28,7 +28,7 @@ public class StartAt : Trigger public StartAt(Input dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => DateTime = dateTime; /// - public StartAt(Func dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) + public StartAt(Func dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input(dateTime), source, line) { } @@ -41,29 +41,30 @@ public class StartAt : Trigger } /// - public StartAt(Func> dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) + public StartAt(Func> dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input(dateTime), source, line) { } /// - public StartAt(Func dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) + public StartAt(Func dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input(dateTime), source, line) { } /// - public StartAt(DateTimeOffset dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => + public StartAt(DateTimeOffset dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => DateTime = new Input(dateTime); /// - public StartAt(Variable dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => + public StartAt(Variable dateTime, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) => DateTime = new Input(dateTime); /// /// The timestamp at which the workflow should be triggered. /// - [Input] public Input DateTime { get; set; } = default!; + [Input] + public Input DateTime { get; set; } = default!; /// protected override object GetTriggerPayload(TriggerIndexingContext context) @@ -73,27 +74,29 @@ public class StartAt : Trigger } /// - protected override void Execute(ActivityExecutionContext context) + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { - // If external input was received, it means this activity got triggered and does not need to create a bookmark. - if (context.TryGetWorkflowInput(InputKey, out _)) + if (context.IsTriggerOfWorkflow()) + { + await context.CompleteActivityAsync(); return; - - // No external input received, so create a bookmark. + } + var executeAt = context.ExpressionExecutionContext.Get(DateTime); var clock = context.ExpressionExecutionContext.GetRequiredService(); var now = clock.UtcNow; var logger = context.GetRequiredService>(); + context.JournalData.Add("Executed At", now); + if (executeAt <= now) { - logger.LogDebug("Scheduled trigger time lies in the past ('{Delta}'). Skipping scheduling", now - executeAt); - context.JournalData.Add("Executed At", now); + logger.LogDebug("Scheduled trigger time lies in the past ('{Delta}'). Completing immediately", now - executeAt); + await context.CompleteActivityAsync(); return; } var payload = new StartAtPayload(executeAt); - context.CreateBookmark(payload); } diff --git a/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs b/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs index 0389e2f78..d63098bca 100644 --- a/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs +++ b/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs @@ -1,3 +1,4 @@ +using Elsa.Common.Contracts; using Elsa.Common.Models; using Elsa.Extensions; using Elsa.Scheduling.Activities; @@ -12,20 +13,9 @@ namespace Elsa.Scheduling.Services; /// /// A default implementation of that schedules triggers using . /// -public class DefaultTriggerScheduler : ITriggerScheduler +public class DefaultTriggerScheduler(IWorkflowScheduler workflowScheduler, ISystemClock systemClock, ILogger logger) + : ITriggerScheduler { - private readonly IWorkflowScheduler _workflowScheduler; - private readonly ILogger _logger; - - /// - /// Initializes a new instance of the class. - /// - public DefaultTriggerScheduler(IWorkflowScheduler workflowScheduler, ILogger logger) - { - _workflowScheduler = workflowScheduler; - _logger = logger; - } - /// public async Task ScheduleAsync(IEnumerable triggers, CancellationToken cancellationToken = default) { @@ -34,6 +24,7 @@ public class DefaultTriggerScheduler : ITriggerScheduler var timerTriggers = triggerList.Filter(); var startAtTriggers = triggerList.Filter(); var cronTriggers = triggerList.Filter(); + var now = systemClock.UtcNow; // Schedule each Timer trigger. foreach (var trigger in timerTriggers) @@ -47,13 +38,21 @@ public class DefaultTriggerScheduler : ITriggerScheduler TriggerActivityId = trigger.ActivityId, Input = input }; - await _workflowScheduler.ScheduleRecurringAsync(trigger.Id, request, startAt, interval, cancellationToken); + await workflowScheduler.ScheduleRecurringAsync(trigger.Id, request, startAt, interval, cancellationToken); } // Schedule each StartAt trigger. 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; + } + var input = new { ExecuteAt = executeAt }.ToDictionary(); var request = new DispatchWorkflowDefinitionRequest { @@ -63,7 +62,7 @@ public class DefaultTriggerScheduler : ITriggerScheduler Input = input }; - await _workflowScheduler.ScheduleAtAsync(trigger.Id, request, executeAt, cancellationToken); + await workflowScheduler.ScheduleAtAsync(trigger.Id, request, executeAt, cancellationToken); } // Schedule each Cron trigger. @@ -81,11 +80,11 @@ public class DefaultTriggerScheduler : ITriggerScheduler }; try { - await _workflowScheduler.ScheduleCronAsync(trigger.Id, request, cronExpression, cancellationToken); + await workflowScheduler.ScheduleCronAsync(trigger.Id, request, cronExpression, cancellationToken); } catch (FormatException ex) { - _logger.LogWarning($"Cron expression format error: {ex.Message}. CronExpression: {cronExpression}"); + logger.LogWarning($"Cron expression format error: {ex.Message}. CronExpression: {cronExpression}"); } } } @@ -109,6 +108,6 @@ public class DefaultTriggerScheduler : ITriggerScheduler // Unschedule each trigger. foreach (var trigger in filteredTriggers) - await _workflowScheduler.UnscheduleAsync(trigger.Id, cancellationToken); + await workflowScheduler.UnscheduleAsync(trigger.Id, cancellationToken); } }