using Elsa.Common;
using Elsa.Extensions;
using Elsa.Scheduling.Activities;
using Elsa.Scheduling.Bookmarks;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Entities;
using Microsoft.Extensions.Logging;
namespace Elsa.Scheduling.Services;
///
/// A default implementation of that schedules triggers using .
///
public class DefaultTriggerScheduler(IWorkflowScheduler workflowScheduler, ISystemClock systemClock, ILogger logger)
: ITriggerScheduler
{
///
public async Task ScheduleAsync(IEnumerable triggers, CancellationToken cancellationToken = default)
{
var triggerList = triggers.ToList();
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)
{
var (startAt, interval) = trigger.GetPayload();
var input = new { StartAt = startAt, Interval = interval }.ToDictionary();
var request = new ScheduleNewWorkflowInstanceRequest
{
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionVersionId(trigger.WorkflowDefinitionVersionId),
TriggerActivityId = trigger.ActivityId,
Input = input
};
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 ScheduleNewWorkflowInstanceRequest
{
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionVersionId(trigger.WorkflowDefinitionVersionId),
TriggerActivityId = trigger.ActivityId,
Input = input
};
await workflowScheduler.ScheduleAtAsync(trigger.Id, request, executeAt, cancellationToken);
}
// Schedule each Cron trigger.
foreach (var trigger in cronTriggers)
{
var payload = trigger.GetPayload();
var cronExpression = payload.CronExpression;
if (string.IsNullOrWhiteSpace(cronExpression))
{
logger.LogWarning("Cron expression is empty. TriggerId: {TriggerId}. Skipping scheduling of this trigger", trigger.Id);
continue;
}
var input = new { CronExpression = cronExpression }.ToDictionary();
var request = new ScheduleNewWorkflowInstanceRequest
{
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionVersionId(trigger.WorkflowDefinitionVersionId),
TriggerActivityId = trigger.ActivityId,
Input = input
};
try
{
await workflowScheduler.ScheduleCronAsync(trigger.Id, request, cronExpression, cancellationToken);
}
catch (FormatException ex)
{
logger.LogWarning(ex, "Cron expression format error. CronExpression: {CronExpression}", cronExpression);
}
}
}
///
public async Task UnscheduleAsync(IEnumerable triggers, CancellationToken cancellationToken = default)
{
var triggerList = triggers.ToList();
var timerTriggers = triggerList.Filter();
var startAtTriggers = triggerList.Filter();
var cronTriggers = triggerList.Filter();
var filteredTriggers = timerTriggers.Concat(startAtTriggers).Concat(cronTriggers);
foreach (var trigger in filteredTriggers)
await workflowScheduler.UnscheduleAsync(trigger.Id, cancellationToken);
}
}