Implement timer trigger scheduling

This commit is contained in:
Sipke Schoorstra 2022-01-18 11:23:58 +01:00
parent 783f05ffdc
commit dcd70aa192
9 changed files with 124 additions and 44 deletions

View file

@ -5,6 +5,6 @@ namespace Elsa.Activities.Scheduling.Contracts;
public interface IJobScheduler
{
Task ScheduleAsync(IJob job, ISchedule schedule, CancellationToken cancellationToken = default);
Task ClearAsync(CancellationToken cancellationToken = default);
Task ScheduleAsync(IJob job, ISchedule schedule, string[]? groupKeys = default, CancellationToken cancellationToken = default);
Task ClearAsync(string[]? groupKeys = default, CancellationToken cancellationToken = default);
}

View file

@ -0,0 +1,14 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Persistence.Entities;
namespace Elsa.Activities.Scheduling.Contracts;
/// <summary>
/// Schedules jobs for the specified list of workflow triggers.
/// </summary>
public interface IWorkflowTriggerScheduler
{
Task ScheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default);
}

View file

@ -1,5 +1,6 @@
using Elsa.Activities.Scheduling.Contracts;
using Elsa.Activities.Scheduling.Handlers;
using Elsa.Activities.Scheduling.HostedServices;
using Elsa.Activities.Scheduling.Jobs;
using Elsa.Activities.Scheduling.Services;
using Elsa.Mediator.Extensions;
@ -15,7 +16,9 @@ public static class ServiceCollectionExtensions
.AddSingleton<IJobManager, JobManager>()
.AddSingleton<IJobHandler, RunWorkflowJobHandler>()
.AddSingleton<IJobHandler, ResumeWorkflowJobHandler>()
.AddNotificationHandlersFrom<ScheduleWorkflowsHandler>();
.AddSingleton<IWorkflowTriggerScheduler, WorkflowTriggerScheduler>()
.AddNotificationHandlersFrom<ScheduleWorkflowsHandler>()
.AddHostedService<ScheduleWorkflowsHostedService>();
serviceProvider.ConfigureServices(services);
return services;

View file

@ -1,40 +1,15 @@
using System.Linq;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.Scheduling.Contracts;
using Elsa.Activities.Scheduling.Jobs;
using Elsa.Activities.Scheduling.Schedules;
using Elsa.Mediator.Contracts;
using Elsa.Persistence.Extensions;
using Elsa.Persistence.Requests;
using Elsa.Runtime.Notifications;
namespace Elsa.Activities.Scheduling.Handlers;
// Updates scheduled jobs based on the updated workflow triggers.
public class ScheduleWorkflowsHandler : INotificationHandler<TriggerIndexingFinished>
{
private readonly IRequestSender _requestSender;
private readonly IJobScheduler _jobScheduler;
public ScheduleWorkflowsHandler(IRequestSender requestSender, IJobScheduler jobScheduler)
{
_requestSender = requestSender;
_jobScheduler = jobScheduler;
}
public async Task HandleAsync(TriggerIndexingFinished notification, CancellationToken cancellationToken)
{
// Unschedule everything.
await _jobScheduler.ClearAsync(cancellationToken);
// Select all Timer triggers.
var timerTriggers = notification.Triggers.Filter<Timer>().ToList();
foreach (var trigger in timerTriggers)
{
var (dateTime, timeSpan) = JsonSerializer.Deserialize<TimerPayload>(trigger.Payload!)!;
await _jobScheduler.ScheduleAsync(new RunWorkflowJob(trigger.WorkflowDefinitionId), new RecurringSchedule(dateTime, timeSpan), cancellationToken);
}
}
private readonly IWorkflowTriggerScheduler _workflowTriggerScheduler;
public ScheduleWorkflowsHandler(IWorkflowTriggerScheduler workflowTriggerScheduler) => _workflowTriggerScheduler = workflowTriggerScheduler;
public async Task HandleAsync(TriggerIndexingFinished notification, CancellationToken cancellationToken) => await _workflowTriggerScheduler.ScheduleTriggersAsync(notification.Triggers, cancellationToken);
}

View file

@ -0,0 +1,30 @@
using System.Collections.Immutable;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.Scheduling.Contracts;
using Elsa.Mediator.Contracts;
using Elsa.Persistence.Requests;
using Microsoft.Extensions.Hosting;
namespace Elsa.Activities.Scheduling.HostedServices;
/// <summary>
/// Loads all timer-specific workflow triggers from the database and create scheduled jobs for them.
/// </summary>
public class ScheduleWorkflowsHostedService : BackgroundService
{
private readonly IRequestSender _requestSender;
private readonly IWorkflowTriggerScheduler _workflowTriggerScheduler;
public ScheduleWorkflowsHostedService(IRequestSender requestSender, IWorkflowTriggerScheduler workflowTriggerScheduler)
{
_requestSender = requestSender;
_workflowTriggerScheduler = workflowTriggerScheduler;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
var timerTriggers = (await _requestSender.RequestAsync(FindWorkflowTriggers.ForTrigger<Timer>(), stoppingToken)).ToImmutableList();
await _workflowTriggerScheduler.ScheduleTriggersAsync(timerTriggers, stoppingToken);
}
}

View file

@ -0,0 +1,49 @@
using System.Collections.Generic;
using System.Linq;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.Scheduling.Contracts;
using Elsa.Activities.Scheduling.Jobs;
using Elsa.Activities.Scheduling.Schedules;
using Elsa.Persistence.Entities;
using Elsa.Persistence.Extensions;
namespace Elsa.Activities.Scheduling.Services;
public class WorkflowTriggerScheduler : IWorkflowTriggerScheduler
{
private const string RootGroupKey = "WorkflowDefinition";
private readonly IJobScheduler _jobScheduler;
public WorkflowTriggerScheduler(IJobScheduler jobScheduler)
{
_jobScheduler = jobScheduler;
}
public async Task ScheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default)
{
var triggerList = triggers.ToList();
// Select all Timer triggers.
var timerTriggers = triggerList.Filter<Timer>().ToList();
// 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);
}
foreach (var trigger in timerTriggers)
{
// Schedule trigger.
var (dateTime, timeSpan) = JsonSerializer.Deserialize<TimerPayload>(trigger.Payload!)!;
var groupKeys = new[] { RootGroupKey, trigger.WorkflowDefinitionId };
await _jobScheduler.ScheduleAsync(new RunWorkflowJob(trigger.WorkflowDefinitionId), new RecurringSchedule(dateTime, timeSpan), groupKeys, cancellationToken);
}
}
}

View file

@ -13,7 +13,7 @@ namespace Elsa.Modules.Quartz.Services;
public class QuartzJobScheduler : IElsaJobScheduler
{
public const string JobDataKey = "ElsaJob";
public const string GroupKey = "ElsaJobs";
public const string RootGroupKey = "ElsaJobs";
private readonly IElsaJobSerializer _elsaJobSerializer;
private readonly ISchedulerFactory _schedulerFactory;
private readonly ILogger _logger;
@ -25,23 +25,24 @@ public class QuartzJobScheduler : IElsaJobScheduler
_logger = logger;
}
public async Task ScheduleAsync(IElsaJob job, IElsaSchedule schedule, CancellationToken cancellationToken = default)
public async Task ScheduleAsync(IElsaJob job, IElsaSchedule schedule, string[]? groupKeys, CancellationToken cancellationToken = default)
{
var quartzTrigger = CreateTrigger(job, schedule);
var quartzTrigger = CreateTrigger(job, schedule, groupKeys);
await ScheduleJob(quartzTrigger, cancellationToken);
}
public async Task ClearAsync(CancellationToken cancellationToken = default)
public async Task ClearAsync(string[]? groupKeys = default, CancellationToken cancellationToken = default)
{
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
var jobKeys = await scheduler.GetJobKeys(GroupMatcher<JobKey>.GroupEquals(GroupKey), cancellationToken);
var groupKey = BuildGroupKey(groupKeys);
var jobKeys = await scheduler.GetJobKeys(GroupMatcher<JobKey>.GroupStartsWith(groupKey), cancellationToken);
await scheduler.DeleteJobs(jobKeys, cancellationToken);
}
private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken)
{
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
try
{
await scheduler.UnscheduleJob(trigger.Key, cancellationToken);
@ -52,11 +53,12 @@ public class QuartzJobScheduler : IElsaJobScheduler
_logger.LogWarning(e, "Failed to schedule trigger {TriggerKey}", trigger.Key.ToString());
}
}
private ITrigger CreateTrigger(IElsaJob job, IElsaSchedule schedule)
private ITrigger CreateTrigger(IElsaJob job, IElsaSchedule schedule, string[]? groupKeys)
{
var jobName = job.GetType().Name;
var triggerKey = new TriggerKey(job.JobId, GroupKey);
var groupKey = BuildGroupKey(groupKeys);
var triggerKey = new TriggerKey(job.JobId, groupKey);
var json = _elsaJobSerializer.Serialize(job);
var builder = TriggerBuilder.Create().ForJob(jobName).WithIdentity(triggerKey).UsingJobData(JobDataKey, json);
@ -85,4 +87,11 @@ public class QuartzJobScheduler : IElsaJobScheduler
return builder.Build();
}
private static string BuildGroupKey(string[]? groupKeys)
{
var groupKeyInputs = new[] { RootGroupKey }.Concat(groupKeys ?? Array.Empty<string>());
return string.Join(":", groupKeyInputs);
}
}

View file

@ -12,8 +12,8 @@ public class ReplaceWorkflowTriggersHandler : ICommandHandler<ReplaceWorkflowTri
public async Task<Unit> HandleAsync(ReplaceWorkflowTriggers command, CancellationToken cancellationToken)
{
var ids = command.WorkflowTriggers.Select(x => x.Id).ToList();
await _store.DeleteWhereAsync(x => ids.Contains(x.Id), cancellationToken);
var workflowDefinitionIds = command.WorkflowTriggers.Select(x => x.WorkflowDefinitionId).Distinct().ToList();
await _store.DeleteWhereAsync(x => workflowDefinitionIds.Contains(x.WorkflowDefinitionId), cancellationToken);
await _store.SaveManyAsync(command.WorkflowTriggers, cancellationToken);
return Unit.Instance;

View file

@ -4,7 +4,7 @@
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information",
"Microsoft.EntityFrameworkCore.Database.Command": "Information"
"Microsoft.EntityFrameworkCore.Database.Command": "Warning"
}
},
"AllowedHosts": "*"