elsa-core/src/modules/Elsa.Modules.Scheduling/Services/WorkflowTriggerScheduler.cs

57 lines
2.1 KiB
C#
Raw Normal View History

2022-01-18 10:23:58 +00:00
using System.Collections.Generic;
using System.Linq;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
2022-02-01 19:45:11 +00:00
using Elsa.Modules.Scheduling.Contracts;
using Elsa.Modules.Scheduling.Jobs;
using Elsa.Modules.Scheduling.Triggers;
2022-01-18 10:23:58 +00:00
using Elsa.Persistence.Entities;
using Elsa.Persistence.Extensions;
2022-03-07 21:56:16 +00:00
using Elsa.Jobs.Contracts;
using Elsa.Jobs.Schedules;
2022-02-01 19:45:11 +00:00
using Timer = Elsa.Modules.Scheduling.Triggers.Timer;
2022-01-18 10:23:58 +00:00
2022-02-01 19:45:11 +00:00
namespace Elsa.Modules.Scheduling.Services;
2022-01-18 10:23:58 +00:00
public class WorkflowTriggerScheduler : IWorkflowTriggerScheduler
{
private const string RootGroupKey = "WorkflowDefinition";
2022-02-22 11:51:59 +00:00
2022-01-18 10:23:58 +00:00
private readonly IJobScheduler _jobScheduler;
public WorkflowTriggerScheduler(IJobScheduler jobScheduler)
{
_jobScheduler = jobScheduler;
}
2022-02-22 11:51:59 +00:00
2022-01-18 10:23:58 +00:00
public async Task ScheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default)
{
// Select all Timer triggers.
2022-02-22 11:51:59 +00:00
var timerTriggers = triggers.Filter<Timer>().ToList();
// Schedule each trigger.
foreach (var trigger in timerTriggers)
{
// Schedule trigger.
2022-03-07 11:01:27 +00:00
var (dateTime, timeSpan) = JsonSerializer.Deserialize<TimerPayload>(trigger.Data!)!;
2022-02-22 11:51:59 +00:00
var groupKeys = new[] { RootGroupKey, trigger.WorkflowDefinitionId };
2022-03-15 12:01:10 +00:00
await _jobScheduler.ScheduleAsync(new RunWorkflowJob(trigger.WorkflowDefinitionId), trigger.WorkflowDefinitionId, new RecurringSchedule(dateTime, timeSpan), groupKeys, cancellationToken);
2022-02-22 11:51:59 +00:00
}
}
public async Task UnscheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default)
{
// Select all Timer triggers.
var timerTriggers = triggers.Filter<Timer>().ToList();
2022-01-18 10:23:58 +00:00
// Unschedule all triggers for the distinct set of affected workflows.
var workflowDefinitionIds = timerTriggers.Select(x => x.WorkflowDefinitionId).Distinct().ToList();
foreach (var workflowDefinitionId in workflowDefinitionIds)
{
var groupKeys = new[] { RootGroupKey, workflowDefinitionId };
await _jobScheduler.ClearAsync(groupKeys, cancellationToken);
}
}
}