diff --git a/src/activities/Elsa.Activities.Scheduling/Contracts/IJobScheduler.cs b/src/activities/Elsa.Activities.Scheduling/Contracts/IJobScheduler.cs index 7d6e791a5..4cc9322e7 100644 --- a/src/activities/Elsa.Activities.Scheduling/Contracts/IJobScheduler.cs +++ b/src/activities/Elsa.Activities.Scheduling/Contracts/IJobScheduler.cs @@ -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); } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Contracts/IWorkflowTriggerScheduler.cs b/src/activities/Elsa.Activities.Scheduling/Contracts/IWorkflowTriggerScheduler.cs new file mode 100644 index 000000000..646bfdfce --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Contracts/IWorkflowTriggerScheduler.cs @@ -0,0 +1,14 @@ +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Persistence.Entities; + +namespace Elsa.Activities.Scheduling.Contracts; + +/// +/// Schedules jobs for the specified list of workflow triggers. +/// +public interface IWorkflowTriggerScheduler +{ + Task ScheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Scheduling/Extensions/ServiceCollectionExtensions.cs index 3811cebc0..1da181e3c 100644 --- a/src/activities/Elsa.Activities.Scheduling/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Scheduling/Extensions/ServiceCollectionExtensions.cs @@ -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() .AddSingleton() .AddSingleton() - .AddNotificationHandlersFrom(); + .AddSingleton() + .AddNotificationHandlersFrom() + .AddHostedService(); serviceProvider.ConfigureServices(services); return services; diff --git a/src/activities/Elsa.Activities.Scheduling/Handlers/ScheduleWorkflowsHandler.cs b/src/activities/Elsa.Activities.Scheduling/Handlers/ScheduleWorkflowsHandler.cs index ac20798b9..77fb38c33 100644 --- a/src/activities/Elsa.Activities.Scheduling/Handlers/ScheduleWorkflowsHandler.cs +++ b/src/activities/Elsa.Activities.Scheduling/Handlers/ScheduleWorkflowsHandler.cs @@ -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 { - 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().ToList(); - - foreach (var trigger in timerTriggers) - { - var (dateTime, timeSpan) = JsonSerializer.Deserialize(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); } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/HostedServices/ScheduleWorkflowsHostedService.cs b/src/activities/Elsa.Activities.Scheduling/HostedServices/ScheduleWorkflowsHostedService.cs new file mode 100644 index 000000000..422855e50 --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/HostedServices/ScheduleWorkflowsHostedService.cs @@ -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; + +/// +/// Loads all timer-specific workflow triggers from the database and create scheduled jobs for them. +/// +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(), stoppingToken)).ToImmutableList(); + await _workflowTriggerScheduler.ScheduleTriggersAsync(timerTriggers, stoppingToken); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Services/WorkflowTriggerScheduler.cs b/src/activities/Elsa.Activities.Scheduling/Services/WorkflowTriggerScheduler.cs new file mode 100644 index 000000000..8ed4d0c49 --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Services/WorkflowTriggerScheduler.cs @@ -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 triggers, CancellationToken cancellationToken = default) + { + var triggerList = triggers.ToList(); + + // Select all Timer triggers. + var timerTriggers = triggerList.Filter().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(trigger.Payload!)!; + var groupKeys = new[] { RootGroupKey, trigger.WorkflowDefinitionId }; + await _jobScheduler.ScheduleAsync(new RunWorkflowJob(trigger.WorkflowDefinitionId), new RecurringSchedule(dateTime, timeSpan), groupKeys, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.Quartz/Services/QuartzJobScheduler.cs b/src/modules/Elsa.Modules.Quartz/Services/QuartzJobScheduler.cs index a64c33c97..cd6f48afc 100644 --- a/src/modules/Elsa.Modules.Quartz/Services/QuartzJobScheduler.cs +++ b/src/modules/Elsa.Modules.Quartz/Services/QuartzJobScheduler.cs @@ -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.GroupEquals(GroupKey), cancellationToken); + var groupKey = BuildGroupKey(groupKeys); + var jobKeys = await scheduler.GetJobKeys(GroupMatcher.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()); + return string.Join(":", groupKeyInputs); + } + } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.EntityFrameworkCore/Handlers/Commands/ReplaceWorkflowTriggersHandler.cs b/src/persistence/Elsa.Persistence.EntityFrameworkCore/Handlers/Commands/ReplaceWorkflowTriggersHandler.cs index f074a283f..e7389a692 100644 --- a/src/persistence/Elsa.Persistence.EntityFrameworkCore/Handlers/Commands/ReplaceWorkflowTriggersHandler.cs +++ b/src/persistence/Elsa.Persistence.EntityFrameworkCore/Handlers/Commands/ReplaceWorkflowTriggersHandler.cs @@ -12,8 +12,8 @@ public class ReplaceWorkflowTriggersHandler : ICommandHandler 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; diff --git a/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json b/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json index 5c64e151c..778fcef3d 100644 --- a/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json +++ b/src/samples/aspnet/Elsa.Samples.Web1/appsettings.json @@ -4,7 +4,7 @@ "Default": "Information", "Microsoft": "Warning", "Microsoft.Hosting.Lifetime": "Information", - "Microsoft.EntityFrameworkCore.Database.Command": "Information" + "Microsoft.EntityFrameworkCore.Database.Command": "Warning" } }, "AllowedHosts": "*"