diff --git a/Elsa.sln b/Elsa.sln index ebe7084ae..be7573b3c 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -141,7 +141,7 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Samples.RebusWorker", EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Samples.RunChildWorkflowWorker", "src\samples\worker\Elsa.Samples.RunChildWorkflowWorker\Elsa.Samples.RunChildWorkflowWorker.csproj", "{A4DF29BD-95CC-4A86-96EC-FB6692270E68}" EndProject -Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Samples.Timers", "src\samples\worker\Elsa.Samples.Timers\Elsa.Samples.Timers.csproj", "{FD2D1BD0-1229-4DCF-BE70-6BFD396DD489}" +Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Samples.Timers.Quartz", "src\samples\worker\Elsa.Samples.Timers\Elsa.Samples.Timers.Quartz.csproj", "{FD2D1BD0-1229-4DCF-BE70-6BFD396DD489}" EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Samples.WhileLoopWorker", "src\samples\worker\Elsa.Samples.WhileLoopWorker\Elsa.Samples.WhileLoopWorker.csproj", "{EA3832C4-8079-4E84-AF8C-12744E4EFD44}" EndProject @@ -167,6 +167,13 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "ElsaDashboard.Backend", "sr EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "ElsaDashboard.WebAssembly", "src\dashboards\blazor\ElsaDashboard.WebAssembly\ElsaDashboard.WebAssembly.csproj", "{0AA2A003-C79E-4B43-803E-E1150254D2B0}" EndProject +Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "timers", "timers", "{37D2E628-5CC1-4C8E-8C02-3B0A29E2ACAA}" +EndProject +Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Activities.Timers.Quartz", "src\activities\Elsa.Activities.Timers.Quartz\Elsa.Activities.Timers.Quartz.csproj", "{08D9CB18-2BFB-4CF2-9763-5EE913F56D3E}" +EndProject +Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Activities.Timers.Hangfire", "src\activities\Elsa.Activities.Timers.Hangfire\Elsa.Activities.Timers.Hangfire.csproj", "{34FB968B-A147-4811-823D-D206FD0027AD}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.Timers.Hangfire", "src\samples\worker\Elsa.Samples.Timers.Hangfire\Elsa.Samples.Timers.Hangfire.csproj", "{7BDF193D-CD69-4D47-B7BF-A8A1D93EBCB9}" Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Activities.Entity", "src\activities\Elsa.Activities.Entity\Elsa.Activities.Entity.csproj", "{A0A8431F-3A2F-433A-BB0D-EB9627DFA739}" EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Samples.EntityChanged", "src\samples\console\Elsa.Samples.EntityChanged\Elsa.Samples.EntityChanged.csproj", "{A3149FE4-FD39-4DC2-AA9F-208C7216BC03}" @@ -403,6 +410,18 @@ Global {0AA2A003-C79E-4B43-803E-E1150254D2B0}.Debug|Any CPU.Build.0 = Debug|Any CPU {0AA2A003-C79E-4B43-803E-E1150254D2B0}.Release|Any CPU.ActiveCfg = Release|Any CPU {0AA2A003-C79E-4B43-803E-E1150254D2B0}.Release|Any CPU.Build.0 = Release|Any CPU + {08D9CB18-2BFB-4CF2-9763-5EE913F56D3E}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {08D9CB18-2BFB-4CF2-9763-5EE913F56D3E}.Debug|Any CPU.Build.0 = Debug|Any CPU + {08D9CB18-2BFB-4CF2-9763-5EE913F56D3E}.Release|Any CPU.ActiveCfg = Release|Any CPU + {08D9CB18-2BFB-4CF2-9763-5EE913F56D3E}.Release|Any CPU.Build.0 = Release|Any CPU + {34FB968B-A147-4811-823D-D206FD0027AD}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {34FB968B-A147-4811-823D-D206FD0027AD}.Debug|Any CPU.Build.0 = Debug|Any CPU + {34FB968B-A147-4811-823D-D206FD0027AD}.Release|Any CPU.ActiveCfg = Release|Any CPU + {34FB968B-A147-4811-823D-D206FD0027AD}.Release|Any CPU.Build.0 = Release|Any CPU + {7BDF193D-CD69-4D47-B7BF-A8A1D93EBCB9}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {7BDF193D-CD69-4D47-B7BF-A8A1D93EBCB9}.Debug|Any CPU.Build.0 = Debug|Any CPU + {7BDF193D-CD69-4D47-B7BF-A8A1D93EBCB9}.Release|Any CPU.ActiveCfg = Release|Any CPU + {7BDF193D-CD69-4D47-B7BF-A8A1D93EBCB9}.Release|Any CPU.Build.0 = Release|Any CPU {A0A8431F-3A2F-433A-BB0D-EB9627DFA739}.Debug|Any CPU.ActiveCfg = Debug|Any CPU {A0A8431F-3A2F-433A-BB0D-EB9627DFA739}.Debug|Any CPU.Build.0 = Debug|Any CPU {A0A8431F-3A2F-433A-BB0D-EB9627DFA739}.Release|Any CPU.ActiveCfg = Release|Any CPU @@ -435,7 +454,7 @@ Global {1D63E1B4-2386-4BEC-9090-C3F3CBF25703} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} {B43B546E-23F3-46E8-ACB7-D04F05CDA180} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925} {D20FCB88-9DCA-49EA-9CC2-5B94DD935D30} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} - {E4B71DC4-3E73-49C3-9B4E-CA13909F222D} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} + {E4B71DC4-3E73-49C3-9B4E-CA13909F222D} = {37D2E628-5CC1-4C8E-8C02-3B0A29E2ACAA} {C4939482-9447-47D1-B6AF-E9F8C6323841} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925} {823F80B6-E241-43F7-83BD-29FBB58A77F9} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} @@ -500,6 +519,10 @@ Global {CCF8D59E-3672-4238-A900-CE879ABF02D2} = {E02BD53A-5B08-47A0-920E-DD537A6AFAAA} {B5EAB378-0002-47FD-A0C1-8931DB577006} = {E02BD53A-5B08-47A0-920E-DD537A6AFAAA} {0AA2A003-C79E-4B43-803E-E1150254D2B0} = {E02BD53A-5B08-47A0-920E-DD537A6AFAAA} + {37D2E628-5CC1-4C8E-8C02-3B0A29E2ACAA} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} + {08D9CB18-2BFB-4CF2-9763-5EE913F56D3E} = {37D2E628-5CC1-4C8E-8C02-3B0A29E2ACAA} + {34FB968B-A147-4811-823D-D206FD0027AD} = {37D2E628-5CC1-4C8E-8C02-3B0A29E2ACAA} + {7BDF193D-CD69-4D47-B7BF-A8A1D93EBCB9} = {E42743A0-FBDD-4150-9D53-6000496D9B87} {A0A8431F-3A2F-433A-BB0D-EB9627DFA739} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} {A3149FE4-FD39-4DC2-AA9F-208C7216BC03} = {FC9F520F-BA51-4AD2-BFEE-EF787798E734} {C3842132-35BC-4F6D-85CC-7098887F4307} = {7CD5C8D5-EC78-4A99-A514-F01CA7197AC8} diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/Elsa.Activities.Timers.Hangfire.csproj b/src/activities/Elsa.Activities.Timers.Hangfire/Elsa.Activities.Timers.Hangfire.csproj new file mode 100644 index 000000000..39b013bb9 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/Elsa.Activities.Timers.Hangfire.csproj @@ -0,0 +1,26 @@ + + + + + + + netstandard2.0 + + Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application. + This package provides Hangfire timer provider. + + + elsa, workflows, timers, background tasks + + + + + + + + + + + + + diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/Extensions/IBackgroundJobClientExtensions.cs b/src/activities/Elsa.Activities.Timers.Hangfire/Extensions/IBackgroundJobClientExtensions.cs new file mode 100644 index 000000000..6ad21792f --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/Extensions/IBackgroundJobClientExtensions.cs @@ -0,0 +1,34 @@ +using System; +using System.Linq; + +using Elsa.Activities.Timers.Hangfire.Jobs; +using Elsa.Activities.Timers.Hangfire.Models; + +using Hangfire; + +namespace Elsa.Activities.Timers +{ + public static class IBackgroundJobClientExtensions + { + public static void ScheduleWorkflow(this IBackgroundJobClient backgroundJobClient, RunHangfireWorkflowJobModel data, DateTimeOffset dateTimeOffset) + { + backgroundJobClient.UnscheduleJobWhenAlreadyExists(data); + backgroundJobClient.Schedule(job => job.ExecuteAsync(data), dateTimeOffset); + } + + public static void UnscheduleJobWhenAlreadyExists(this IBackgroundJobClient backgroundJobClient, RunHangfireWorkflowJobModel data) + { + var identity = data.GetIdentity(); + var monitor = JobStorage.Current.GetMonitoringApi(); + var workflowJobType = typeof(RunHangfireWorkflowJob); + + var jobs = monitor.ScheduledJobs(0, int.MaxValue) + .Where(x => x.Value.Job.Type == workflowJobType && ((RunHangfireWorkflowJobModel)x.Value.Job.Args[0]).GetIdentity() == identity); + + foreach (var job in jobs) + { + BackgroundJob.Delete(job.Key); + } + } + } +} diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/Extensions/TimersOptionsExtensions.cs b/src/activities/Elsa.Activities.Timers.Hangfire/Extensions/TimersOptionsExtensions.cs new file mode 100644 index 000000000..b55991498 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/Extensions/TimersOptionsExtensions.cs @@ -0,0 +1,45 @@ +using System; + +using Elsa.Activities.Timers.Hangfire.Services; +using Elsa.Activities.Timers.Options; +using Elsa.Activities.Timers.Services; + +using Hangfire; + +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Activities.Timers.Hangfire.Extensions +{ + public static class TimersOptionsExtensions + { + /// + /// Add Elsa Hangfire Services for background processing + /// + /// + public static void UseHangfire(this TimersOptions timersOptions) + { + timersOptions.Services + .AddSingleton() + .AddSingleton(); + } + + /// + /// Add Elsa Hangfire Services for background processing and Hangfire Services + /// + /// + /// Only if Hangfire is not already registered in DI + /// + /// + /// Hangfire settings + public static void UseHangfire(this TimersOptions timersOptions, Action configure) + { + timersOptions.UseHangfire(); + + // Add Hangfire services. + timersOptions.Services.AddHangfire(configure); + + // Add the processing server as IHostedService + timersOptions.Services.AddHangfireServer(); + } + } +} diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/FodyWeavers.xml b/src/activities/Elsa.Activities.Timers.Hangfire/FodyWeavers.xml new file mode 100644 index 000000000..4c4041f8a --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/FodyWeavers.xml @@ -0,0 +1,3 @@ + + + \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/Jobs/RunHangfireWorkflowJob.cs b/src/activities/Elsa.Activities.Timers.Hangfire/Jobs/RunHangfireWorkflowJob.cs new file mode 100644 index 000000000..06832ef7e --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/Jobs/RunHangfireWorkflowJob.cs @@ -0,0 +1,92 @@ +using System; +using System.Threading.Tasks; + +using Elsa.Activities.Timers.Hangfire.Models; +using Elsa.Activities.Timers.Services; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Services; + +using Hangfire; +using Hangfire.States; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; + +namespace Elsa.Activities.Timers.Hangfire.Jobs +{ + public class RunHangfireWorkflowJob + { + private readonly IWorkflowRunner _workflowRunner; + private readonly IWorkflowRegistry _workflowRegistry; + private readonly IWorkflowInstanceStore _workflowInstanceManager; + private readonly IServiceProvider _serviceProvider; + + public RunHangfireWorkflowJob(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, IServiceProvider serviceProvider) + { + _workflowRunner = workflowRunner; + _workflowRegistry = workflowRegistry; + _workflowInstanceManager = workflowInstanceStore; + _serviceProvider = serviceProvider; + } + + public async Task ExecuteAsync(RunHangfireWorkflowJobModel data) + { + var workflowBlueprint = (await _workflowRegistry.GetWorkflowAsync(data.WorkflowDefinitionId, data.TenantId, VersionOptions.Published))!; + + if(workflowBlueprint == null) + { + return; + } + + WorkflowInstance? workflowInstance = null; + + if (data.WorkflowInstanceId == null) + { + if (workflowBlueprint.IsSingleton == false || await _workflowInstanceManager.GetWorkflowIsAlreadyExecutingAsync(data.TenantId, data.WorkflowDefinitionId) == false) + { + await _workflowRunner.RunWorkflowAsync(workflowBlueprint, data.ActivityId); + } + } + else + { + workflowInstance = await GetWorkflowInstanceAsync(data.WorkflowInstanceId); + + if (workflowInstance == null) + { + var logger = _serviceProvider.GetRequiredService>(); + logger.LogError("Could not run Workflow instance with ID {WorkflowInstanceId} because it is not in the database", data.WorkflowInstanceId); + + return; + } + + await _workflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance!, data.ActivityId); + } + + // If it is a RecurringJob and the instance is null, the timer activity is a start trigger. + if (data.IsRecurringJob && (workflowInstance == null || workflowInstance.Status is not (WorkflowStatus.Finished or WorkflowStatus.Cancelled))) + { + var backgroundJobClient = _serviceProvider.GetRequiredService(); + var crontabParer = _serviceProvider.GetRequiredService(); + + backgroundJobClient.ScheduleWorkflow(data, crontabParer.GetNextOccurrence(data.CronExpression!).ToDateTimeOffset()); + } + } + + private async Task GetWorkflowInstanceAsync(string workflowInstanceId) + { + WorkflowInstance? workflowInstance = null; + + for (var i = 0; i < TimerConsts.MaxRetrayGetWorkflow && workflowInstance == null; i++) + { + workflowInstance = await _workflowInstanceManager.GetByIdAsync(workflowInstanceId); + + if (workflowInstance == null) + { + System.Threading.Thread.Sleep(10000); + } + } + + return workflowInstance; + } + } +} diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/Models/RunHangfireWorkflowJobModel.cs b/src/activities/Elsa.Activities.Timers.Hangfire/Models/RunHangfireWorkflowJobModel.cs new file mode 100644 index 000000000..677276b31 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/Models/RunHangfireWorkflowJobModel.cs @@ -0,0 +1,27 @@ +using System; + +using NodaTime; + +namespace Elsa.Activities.Timers.Hangfire.Models +{ + public class RunHangfireWorkflowJobModel + { + public RunHangfireWorkflowJobModel(string workflowDefinitionId, string activityId, string? workflowInstanceId, string? tenantId, string? cronExpression) + { + WorkflowDefinitionId = workflowDefinitionId; + WorkflowInstanceId = workflowInstanceId; + ActivityId = activityId; + TenantId = tenantId; + CronExpression = cronExpression; + } + + public string WorkflowDefinitionId { get; set; } + public string? WorkflowInstanceId { get; set; } + public string ActivityId { get; set; } + public string? TenantId { get; set; } + public string? CronExpression { get; set; } + public bool IsRecurringJob => string.IsNullOrEmpty(CronExpression) == false; + + public string GetIdentity() => $"Elsa-tenant:{TenantId ?? "default"}-workflow-instance:{WorkflowInstanceId ?? WorkflowDefinitionId}-activity:{ActivityId}"; + } +} diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/Services/HangfireCrontabParser.cs b/src/activities/Elsa.Activities.Timers.Hangfire/Services/HangfireCrontabParser.cs new file mode 100644 index 000000000..5d9519d07 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/Services/HangfireCrontabParser.cs @@ -0,0 +1,20 @@ +using System; + +using Elsa.Activities.Timers.Services; + +using NCrontab; + +using NodaTime; + +namespace Elsa.Activities.Timers.Hangfire.Services +{ + public class HangfireCrontabParser : ICrontabParser + { + public Instant GetNextOccurrence(string cronExpression) + { + var schedule = CrontabSchedule.Parse(cronExpression, new CrontabSchedule.ParseOptions { IncludingSeconds = true }); + + return Instant.FromDateTimeUtc(schedule.GetNextOccurrence(DateTime.Now).ToUniversalTime()); + } + } +} diff --git a/src/activities/Elsa.Activities.Timers.Hangfire/Services/HangfireWorkflowScheduler.cs b/src/activities/Elsa.Activities.Timers.Hangfire/Services/HangfireWorkflowScheduler.cs new file mode 100644 index 000000000..515987a41 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Hangfire/Services/HangfireWorkflowScheduler.cs @@ -0,0 +1,86 @@ +using System; +using System.Threading; +using System.Threading.Tasks; + +using Elsa.Activities.Timers.Hangfire.Models; +using Elsa.Activities.Timers.Services; +using Elsa.Services.Models; + +using Hangfire; + +using NodaTime; + +namespace Elsa.Activities.Timers.Hangfire.Services +{ + public class HangfireWorkflowScheduler : IWorkflowScheduler + { + private readonly IBackgroundJobClient _backgroundJobClient; + private readonly ICrontabParser _crontabParser; + + public HangfireWorkflowScheduler(IBackgroundJobClient backgroundJobClient, ICrontabParser crontabParser) + { + _backgroundJobClient = backgroundJobClient; + _crontabParser = crontabParser; + } + + public Task ScheduleWorkflowAsync(IWorkflowBlueprint workflowBlueprint, string activityId, Instant startAt, Duration interval, CancellationToken cancellationToken = default) + { + var data = CreateData(workflowBlueprint, activityId: activityId, cronExpression: interval.ToCronExpression()); + + _backgroundJobClient.ScheduleWorkflow(data, startAt.ToDateTimeOffset()); + + return Task.CompletedTask; + } + + public Task ScheduleWorkflowAsync(IWorkflowBlueprint workflowBlueprint, string activityId, Instant startAt, CancellationToken cancellationToken = default) + { + var data = CreateData(workflowBlueprint, activityId); + + _backgroundJobClient.ScheduleWorkflow(data, startAt.ToDateTimeOffset()); + + return Task.CompletedTask; + } + + public Task ScheduleWorkflowAsync(IWorkflowBlueprint workflowBlueprint, string activityId, string cronExpression, CancellationToken cancellationToken = default) + { + var data = CreateData(workflowBlueprint, activityId, cronExpression: cronExpression); + var instant = _crontabParser.GetNextOccurrence(cronExpression); + + _backgroundJobClient.ScheduleWorkflow(data, instant.ToDateTimeOffset()); + + return Task.CompletedTask; + } + + public Task ScheduleWorkflowAsync(IWorkflowBlueprint workflowBlueprint, string workflowInstanceId, string activityId, Instant startAt, CancellationToken cancellationToken = default) + { + var data = CreateData(workflowBlueprint, activityId, workflowInstanceId); + + _backgroundJobClient.ScheduleWorkflow(data, startAt.ToDateTimeOffset()); + + return Task.CompletedTask; + } + + public Task UnscheduleWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, string activityId, CancellationToken cancellationToken = default) + { + _backgroundJobClient.UnscheduleJobWhenAlreadyExists( + CreateData(workflowExecutionContext.WorkflowBlueprint, + activityId: activityId, + workflowInstanceId: workflowExecutionContext.WorkflowInstance.WorkflowInstanceId) + ); + + return Task.CompletedTask; + } + + private RunHangfireWorkflowJobModel CreateData(IWorkflowBlueprint workflowBlueprint, string activityId, string? workflowInstanceId = null, string? cronExpression = null) => CreateData(workflowBlueprint.Id,activityId, workflowInstanceId, workflowBlueprint.TenantId, cronExpression); + private RunHangfireWorkflowJobModel CreateData(string workflowDefinitionId, string activityId, string? workflowInstanceId = null, string? tenantId = null, string? cronExpression = null) + { + return new RunHangfireWorkflowJobModel( + workflowDefinitionId: workflowDefinitionId, + activityId: activityId, + workflowInstanceId: workflowInstanceId, + tenantId: tenantId, + cronExpression: cronExpression); + } + + } +} diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Elsa.Activities.Timers.Quartz.csproj b/src/activities/Elsa.Activities.Timers.Quartz/Elsa.Activities.Timers.Quartz.csproj new file mode 100644 index 000000000..9508b204c --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Quartz/Elsa.Activities.Timers.Quartz.csproj @@ -0,0 +1,25 @@ + + + + + + + netstandard2.0 + + Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application. + This package provides Quartz timer provider. + + + + + + + + + + + + + + + diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Extensions/TimersOptionsExtensions.cs b/src/activities/Elsa.Activities.Timers.Quartz/Extensions/TimersOptionsExtensions.cs new file mode 100644 index 000000000..62dce5e4f --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Quartz/Extensions/TimersOptionsExtensions.cs @@ -0,0 +1,55 @@ +using System; + +using Elsa.Activities.Timers.Options; +using Elsa.Activities.Timers.Quartz.Jobs; +using Elsa.Activities.Timers.Quartz.Services; +using Elsa.Activities.Timers.Services; + +using Microsoft.Extensions.DependencyInjection; + +using Quartz; + +namespace Elsa.Activities.Timers +{ + public static class TimersOptionsExtensions + { + /// + /// Add Quartz for background processing + /// + /// + public static void UseQuartzProvider(this TimersOptions timersOptions, Action? configureOptions = default, Action? configureQuartz = default) + { + if (configureOptions != null) + timersOptions.Services.Configure(configureOptions); + else + timersOptions.Services.AddOptions(); + + timersOptions.Services.AddQuartz(configure => ConfigureQuartz(configure, configureQuartz)) + .AddQuartzHostedService(ConfigureQuartzHostedService) + .AddSingleton() + .AddSingleton() + .AddTransient(); + } + + private static void ConfigureQuartzHostedService(QuartzHostedServiceOptions options) + { + options.WaitForJobsToComplete = true; + } + + private static void ConfigureQuartz(IServiceCollectionQuartzConfigurator quartz, Action? configureQuartz) + { + quartz.UseMicrosoftDependencyInjectionScopedJobFactory(options => options.AllowDefaultConstructor = true); + quartz.AddJob(job => job.StoreDurably().WithIdentity(nameof(RunQuartzWorkflowJob))); + + if (configureQuartz != null) + { + configureQuartz(quartz); + } + else + { + quartz.UseSimpleTypeLoader(); + quartz.UseInMemoryStore(); + } + } + } +} diff --git a/src/activities/Elsa.Activities.Timers.Quartz/FodyWeavers.xml b/src/activities/Elsa.Activities.Timers.Quartz/FodyWeavers.xml new file mode 100644 index 000000000..4c4041f8a --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Quartz/FodyWeavers.xml @@ -0,0 +1,3 @@ + + + \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/Jobs/RunWorkflowJob.cs b/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs similarity index 66% rename from src/activities/Elsa.Activities.Timers/Jobs/RunWorkflowJob.cs rename to src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs index 5a78f6cb3..dcc10c2b4 100644 --- a/src/activities/Elsa.Activities.Timers/Jobs/RunWorkflowJob.cs +++ b/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs @@ -1,4 +1,6 @@ -using System.Threading.Tasks; +using System.Threading; +using System.Threading.Tasks; + using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications; @@ -6,16 +8,16 @@ using Elsa.Services; using Microsoft.Extensions.Logging; using Quartz; -namespace Elsa.Activities.Timers.Jobs +namespace Elsa.Activities.Timers.Quartz.Jobs { - public class RunWorkflowJob : IJob + public class RunQuartzWorkflowJob : IJob { private readonly IWorkflowRunner _workflowRunner; private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowInstanceStore _workflowInstanceManager; private readonly ILogger _logger; - public RunWorkflowJob(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, ILogger logger) + public RunQuartzWorkflowJob(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, ILogger logger) { _workflowRunner = workflowRunner; _workflowRegistry = workflowRegistry; @@ -51,16 +53,29 @@ namespace Elsa.Activities.Timers.Jobs if (workflowInstance == null) { - _logger.LogWarning("Could not run Workflow instance with ID {WorkflowInstanceId} because it appears not yet to be persisted in the database. Rescheduling.", workflowInstanceId); - var trigger = context.Trigger; - await context.Scheduler.UnscheduleJob(trigger.Key, cancellationToken); - var newTrigger = trigger.GetTriggerBuilder().StartAt(trigger.StartTimeUtc.AddSeconds(10)).Build(); - await context.Scheduler.ScheduleJob(newTrigger, cancellationToken); + _logger.LogError("Could not run Workflow instance with ID {WorkflowInstanceId} because it is not in the database", data.WorkflowInstanceId); return; } await _workflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance!, activityId, cancellationToken: cancellationToken); } } + + private async Task GetWorkflowInstanceAsync(string workflowInstanceId, CancellationToken cancellationToken) + { + WorkflowInstance? workflowInstance = null; + + for (var i = 0; i < TimerConsts.MaxRetrayGetWorkflow && workflowInstance == null; i++) + { + workflowInstance = await _workflowInstanceManager.GetByIdAsync(workflowInstanceId, cancellationToken); + + if (workflowInstance == null) + { + System.Threading.Thread.Sleep(10000); + } + } + + return workflowInstance; + } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzCrontabParser.cs b/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzCrontabParser.cs new file mode 100644 index 000000000..306ce6fb7 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzCrontabParser.cs @@ -0,0 +1,25 @@ +using Elsa.Activities.Timers.Services; + +using NodaTime; + +using Quartz; + +namespace Elsa.Activities.Timers.Quartz.Services +{ + public class QuartzCrontabParser : ICrontabParser + { + private readonly IClock _clock; + + public QuartzCrontabParser(IClock clock) + { + _clock = clock; + } + + public Instant GetNextOccurrence(string cronExpression) + { + var schedule = new CronExpression(cronExpression); + var now = _clock.GetCurrentInstant(); + return Instant.FromDateTimeOffset(schedule.GetTimeAfter(now.ToDateTimeOffset())!.Value); + } + } +} diff --git a/src/activities/Elsa.Activities.Timers/Services/WorkflowScheduler.cs b/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzWorkflowScheduler.cs similarity index 72% rename from src/activities/Elsa.Activities.Timers/Services/WorkflowScheduler.cs rename to src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzWorkflowScheduler.cs index deda96930..46564f21a 100644 --- a/src/activities/Elsa.Activities.Timers/Services/WorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzWorkflowScheduler.cs @@ -1,20 +1,24 @@ using System; using System.Threading; using System.Threading.Tasks; -using Elsa.Activities.Timers.Jobs; + +using Elsa.Activities.Timers.Quartz.Jobs; +using Elsa.Activities.Timers.Services; using Elsa.Services.Models; + using NodaTime; + using Quartz; -namespace Elsa.Activities.Timers.Services +namespace Elsa.Activities.Timers.Quartz.Services { - public class WorkflowScheduler : IWorkflowScheduler + public class QuartzWorkflowScheduler: IWorkflowScheduler { - private static readonly string RunWorkflowJobKey = nameof(RunWorkflowJob); + private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowJob); private readonly ISchedulerFactory _schedulerFactory; private readonly SemaphoreSlim _semaphore = new(1); - public WorkflowScheduler(ISchedulerFactory schedulerFactory) + public QuartzWorkflowScheduler(ISchedulerFactory schedulerFactory) { _schedulerFactory = schedulerFactory; } @@ -45,10 +49,24 @@ namespace Elsa.Activities.Timers.Services { var trigger = CreateTrigger(workflowBlueprint, activityId, workflowInstanceId) .StartAt(startAt.ToDateTimeOffset()).Build(); - + await ScheduleJob(trigger, cancellationToken); } + public async Task UnscheduleWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, string activityId, CancellationToken cancellationToken = default) + { + var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var trigger = CreateTriggerKey(tenantId: workflowExecutionContext.WorkflowBlueprint.TenantId, + workflowDefinitionId: workflowExecutionContext.WorkflowBlueprint.Id, + workflowInstanceId: workflowExecutionContext.WorkflowInstance.WorkflowInstanceId, + activityId: activityId); + + var existingTrigger = await scheduler.GetTrigger(trigger, cancellationToken); + + if (existingTrigger != null) + await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken); + } + private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) { await _semaphore.WaitAsync(cancellationToken); @@ -74,15 +92,19 @@ namespace Elsa.Activities.Timers.Services private TriggerBuilder CreateTrigger(string? tenantId, string workflowDefinitionId, string? workflowInstanceId, string activityId) { - var groupName = $"tenant:{tenantId ?? "default"}-workflow-instance:{workflowInstanceId ?? workflowDefinitionId}"; - return TriggerBuilder.Create() .ForJob(RunWorkflowJobKey) - .WithIdentity($"activity:{activityId}", groupName) + .WithIdentity(CreateTriggerKey(tenantId, workflowDefinitionId, workflowInstanceId, activityId)) .UsingJobData("TenantId", tenantId!) .UsingJobData("WorkflowDefinitionId", workflowDefinitionId) .UsingJobData("WorkflowInstanceId", workflowInstanceId!) .UsingJobData("ActivityId", activityId); } + + private TriggerKey CreateTriggerKey(string? tenantId, string workflowDefinitionId, string? workflowInstanceId, string activityId) + { + var groupName = $"tenant:{tenantId ?? "default"}-workflow-instance:{workflowInstanceId ?? workflowDefinitionId}"; + return new TriggerKey($"activity:{activityId}", groupName); + } } -} \ No newline at end of file +} diff --git a/src/activities/Elsa.Activities.Timers/Activities/ClearTimer/CancelTimer.cs b/src/activities/Elsa.Activities.Timers/Activities/ClearTimer/CancelTimer.cs new file mode 100644 index 000000000..a11ababb4 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers/Activities/ClearTimer/CancelTimer.cs @@ -0,0 +1,32 @@ +using System.Threading.Tasks; + +using Elsa.Activities.Timers.Services; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Services; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.Timers +{ + + [Trigger(Category = "Timers", Description = "Cancel a timer (Cron, StartAt, Timer) so that it is not executed. ")] + public class ClearTimer : Activity + { + private readonly IWorkflowScheduler _workflowScheduler; + + public ClearTimer(IWorkflowScheduler workflowScheduler) + { + _workflowScheduler = workflowScheduler; + } + + [ActivityProperty(Hint = "The id of the timer (Cron, StartAt, Timer) activity, which is to be cleared")] + public string ActivityId { get; set; } = default!; + + protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) + { + await _workflowScheduler.UnscheduleWorkflowAsync(context.WorkflowExecutionContext, ActivityId); + return Done(); + } + } +} diff --git a/src/activities/Elsa.Activities.Timers/Activities/ClearTimer/CancelTimerExtensions.cs b/src/activities/Elsa.Activities.Timers/Activities/ClearTimer/CancelTimerExtensions.cs new file mode 100644 index 000000000..fc360d03e --- /dev/null +++ b/src/activities/Elsa.Activities.Timers/Activities/ClearTimer/CancelTimerExtensions.cs @@ -0,0 +1,24 @@ +using System; + +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.Timers +{ + public static class ClearTimerExtensions + { + public static IActivityBuilder CancelTimer(this IBuilder builder, + Action>? setup = default) => builder.Then(setup); + + public static IActivityBuilder CancelTimer(this IBuilder builder, + Func activityId) => + builder.CancelTimer(setup => setup.Set(x => x.ActivityId, activityId)); + + public static IActivityBuilder CancelTimer(this IBuilder builder, Func activityId) => + builder.CancelTimer(setup => setup.Set(x => x.ActivityId, activityId)); + + public static IActivityBuilder CancelTimer(this IBuilder builder, string activityId) => + builder.CancelTimer(setup => setup.Set(x => x.ActivityId, activityId)); + } +} diff --git a/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs b/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs index 68f86dc90..a7ce54a27 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs @@ -1,3 +1,4 @@ +using System.Runtime.InteropServices; using System.Threading.Tasks; using Elsa.Activities.Timers.Services; using Elsa.ActivityResults; @@ -5,6 +6,7 @@ using Elsa.Attributes; using Elsa.Persistence; using Elsa.Services; using Elsa.Services.Models; + using NodaTime; // ReSharper disable once CheckNamespace @@ -20,12 +22,14 @@ namespace Elsa.Activities.Timers private readonly IClock _clock; private readonly IWorkflowInstanceStore _workflowInstanceStore; private readonly IWorkflowScheduler _workflowScheduler; + private readonly ICrontabParser _crontabParser; public Cron(IWorkflowInstanceStore workflowInstanceStore, IWorkflowScheduler workflowScheduler, IClock clock) { _clock = clock; _workflowInstanceStore = workflowInstanceStore; _workflowScheduler = workflowScheduler; + _crontabParser = crontabParser; } [ActivityProperty(Hint = "Specify a CRON expression. See https://crontab.guru/ for help.")] @@ -45,7 +49,7 @@ namespace Elsa.Activities.Timers var cancellationToken = context.CancellationToken; var workflowBlueprint = context.WorkflowExecutionContext.WorkflowBlueprint; var workflowInstance = context.WorkflowExecutionContext.WorkflowInstance; - var executeAt = GetNextOccurrence(CronExpression); + var executeAt = _crontabParser.GetNextOccurrence(CronExpression); ExecuteAt = executeAt; @@ -59,12 +63,5 @@ namespace Elsa.Activities.Timers } protected override IActivityExecutionResult OnResume() => Done(); - - private Instant GetNextOccurrence(string cronExpression) - { - var schedule = new Quartz.CronExpression(cronExpression); - var now = _clock.GetCurrentInstant(); - return Instant.FromDateTimeOffset(schedule.GetTimeAfter(now.ToDateTimeOffset())!.Value); - } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/Elsa.Activities.Timers.csproj b/src/activities/Elsa.Activities.Timers/Elsa.Activities.Timers.csproj index bb96088d7..9eb98db33 100644 --- a/src/activities/Elsa.Activities.Timers/Elsa.Activities.Timers.csproj +++ b/src/activities/Elsa.Activities.Timers/Elsa.Activities.Timers.csproj @@ -7,6 +7,11 @@ netstandard2.0 Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application. + + 2020 + https://github.com/elsa-workflows/elsa-core + https://github.com/elsa-workflows/elsa-core + GitHub This package provides the following Timer activities: * CronEvent @@ -20,11 +25,4 @@ - - - - - - - diff --git a/src/activities/Elsa.Activities.Timers/Extensions/DurationExtensions.cs b/src/activities/Elsa.Activities.Timers/Extensions/DurationExtensions.cs new file mode 100644 index 000000000..c283dff1e --- /dev/null +++ b/src/activities/Elsa.Activities.Timers/Extensions/DurationExtensions.cs @@ -0,0 +1,21 @@ +using NodaTime; + +namespace NodaTime +{ + public static class DurationExtensions + { + public static string ToCronExpression(this Duration duration) + { + static string CreateCronComponent(int number) + { + return (number > 0 ? $"*/{number}" : "* "); + } + + var cron = CreateCronComponent(duration.Seconds); + cron += ' ' + CreateCronComponent(duration.Minutes); + cron += ' ' + CreateCronComponent(duration.Hours); + cron += ' ' + CreateCronComponent(duration.Days); + return cron + " * *"; + } + } +} diff --git a/src/activities/Elsa.Activities.Timers/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Timers/Extensions/ServiceCollectionExtensions.cs index 6b089bf08..e3d7fb1ab 100644 --- a/src/activities/Elsa.Activities.Timers/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Timers/Extensions/ServiceCollectionExtensions.cs @@ -1,56 +1,30 @@ using System; + using Elsa.Activities.Timers; using Elsa.Activities.Timers.HostedServices; -using Elsa.Activities.Timers.Jobs; -using Elsa.Activities.Timers.Services; +using Elsa.Activities.Timers.Options; using Elsa.Activities.Timers.Triggers; -using Quartz; // ReSharper disable once CheckNamespace namespace Microsoft.Extensions.DependencyInjection { public static class ServiceCollectionExtensions { - public static IServiceCollection AddTimerActivities(this IServiceCollection services, Action? configureOptions = default, Action? configureQuartz = default) + public static IServiceCollection AddTimerActivities(this IServiceCollection services, + Action? configure = default) { - if (configureOptions != null) - services.Configure(configureOptions); - else - services.AddOptions(); + var options = new TimersOptions(services); + configure?.Invoke(options); - return services - .AddQuartz(configure => ConfigureQuartz(configure, configureQuartz)) - .AddQuartzHostedService(ConfigureQuartzHostedService) - .AddSingleton() - .AddTransient() + return services .AddHostedService() .AddActivity() .AddActivity() .AddActivity() + .AddActivity() .AddTriggerProvider() .AddTriggerProvider() .AddTriggerProvider(); } - - private static void ConfigureQuartzHostedService(QuartzHostedServiceOptions options) - { - options.WaitForJobsToComplete = true; - } - - private static void ConfigureQuartz(IServiceCollectionQuartzConfigurator quartz, Action? configureQuartz) - { - quartz.UseMicrosoftDependencyInjectionScopedJobFactory(options => options.AllowDefaultConstructor = true); - quartz.AddJob(job => job.StoreDurably().WithIdentity(nameof(RunWorkflowJob))); - - if (configureQuartz != null) - { - configureQuartz(quartz); - } - else - { - quartz.UseSimpleTypeLoader(); - quartz.UseInMemoryStore(); - } - } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/Options/TimersOptions.cs b/src/activities/Elsa.Activities.Timers/Options/TimersOptions.cs index 12a738dde..30f4510a1 100644 --- a/src/activities/Elsa.Activities.Timers/Options/TimersOptions.cs +++ b/src/activities/Elsa.Activities.Timers/Options/TimersOptions.cs @@ -1,14 +1,17 @@ +using Microsoft.Extensions.DependencyInjection; + using NodaTime; namespace Elsa.Activities.Timers.Options { public class TimersOptions { - public TimersOptions() + public TimersOptions(IServiceCollection services) { - SweepInterval = Duration.FromMinutes(1); + Services = services; } - - public Duration SweepInterval { get; set; } + + public IServiceCollection Services { get; } + } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/Services/ICrontabParser.cs b/src/activities/Elsa.Activities.Timers/Services/ICrontabParser.cs new file mode 100644 index 000000000..2e9d1a19d --- /dev/null +++ b/src/activities/Elsa.Activities.Timers/Services/ICrontabParser.cs @@ -0,0 +1,16 @@ +using NodaTime; + +namespace Elsa.Activities.Timers.Services +{ + /// + /// The providers can support different formats. Quartz, for example, supports years. + public interface ICrontabParser + { + /// + /// Converts a provider dependent cron string to + /// + /// + /// + Instant GetNextOccurrence(string cronExpression); + } +} diff --git a/src/activities/Elsa.Activities.Timers/Services/IWorkflowScheduler.cs b/src/activities/Elsa.Activities.Timers/Services/IWorkflowScheduler.cs index 428ac0d2d..52972a5c8 100644 --- a/src/activities/Elsa.Activities.Timers/Services/IWorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Timers/Services/IWorkflowScheduler.cs @@ -1,4 +1,4 @@ -using System.Threading; +using System.Threading; using System.Threading.Tasks; using Elsa.Services.Models; using NodaTime; @@ -11,5 +11,6 @@ namespace Elsa.Activities.Timers.Services Task ScheduleWorkflowAsync(IWorkflowBlueprint workflowBlueprint, string activityId, Instant startAt, CancellationToken cancellationToken = default); Task ScheduleWorkflowAsync(IWorkflowBlueprint workflowBlueprint, string activityId, string cronExpression, CancellationToken cancellationToken = default); Task ScheduleWorkflowAsync(IWorkflowBlueprint workflowBlueprint, string workflowInstanceId, string activityId, Instant startAt, CancellationToken cancellationToken = default); + Task UnscheduleWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, string activityId, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/TimerConsts.cs b/src/activities/Elsa.Activities.Timers/TimerConsts.cs new file mode 100644 index 000000000..3b3c19ec5 --- /dev/null +++ b/src/activities/Elsa.Activities.Timers/TimerConsts.cs @@ -0,0 +1,11 @@ +using System; +using System.Collections.Generic; +using System.Text; + +namespace Elsa.Activities.Timers +{ + public static class TimerConsts + { + public const int MaxRetrayGetWorkflow = 3; + } +} diff --git a/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Elsa.Samples.ForkJoinTimerAndSignalHttp.csproj b/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Elsa.Samples.ForkJoinTimerAndSignalHttp.csproj index 57eec47d7..f21d0ed42 100644 --- a/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Elsa.Samples.ForkJoinTimerAndSignalHttp.csproj +++ b/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Elsa.Samples.ForkJoinTimerAndSignalHttp.csproj @@ -6,6 +6,7 @@ + diff --git a/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Startup.cs b/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Startup.cs index 6f224bc5d..116739e05 100644 --- a/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Startup.cs +++ b/src/samples/aspnet/Elsa.Samples.ForkJoinTimerAndSignalHttp/Startup.cs @@ -2,6 +2,7 @@ using Elsa.Samples.ForkJoinTimerAndSignalHttp.BackgroundTasks; using Elsa.Samples.ForkJoinTimerAndSignalHttp.Workflows; using Microsoft.AspNetCore.Builder; using Microsoft.Extensions.DependencyInjection; +using Elsa.Activities.Timers; namespace Elsa.Samples.ForkJoinTimerAndSignalHttp { @@ -15,7 +16,7 @@ namespace Elsa.Samples.ForkJoinTimerAndSignalHttp services .AddElsa() .AddConsoleActivities() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddHostedService>() .AddWorkflow(); } diff --git a/src/samples/aspnet/Elsa.Samples.Interrupts/Elsa.Samples.Interrupts.csproj b/src/samples/aspnet/Elsa.Samples.Interrupts/Elsa.Samples.Interrupts.csproj index 06e0c8e81..d3f06d94a 100644 --- a/src/samples/aspnet/Elsa.Samples.Interrupts/Elsa.Samples.Interrupts.csproj +++ b/src/samples/aspnet/Elsa.Samples.Interrupts/Elsa.Samples.Interrupts.csproj @@ -7,6 +7,7 @@ + diff --git a/src/samples/aspnet/Elsa.Samples.Interrupts/Startup.cs b/src/samples/aspnet/Elsa.Samples.Interrupts/Startup.cs index 19bb61269..aa21c13e7 100644 --- a/src/samples/aspnet/Elsa.Samples.Interrupts/Startup.cs +++ b/src/samples/aspnet/Elsa.Samples.Interrupts/Startup.cs @@ -3,6 +3,7 @@ using Elsa.Samples.Interrupts.Workflows; using Microsoft.AspNetCore.Builder; using Microsoft.AspNetCore.Hosting; using Microsoft.Extensions.DependencyInjection; +using Elsa.Activities.Timers; namespace Elsa.Samples.Interrupts { @@ -16,7 +17,7 @@ namespace Elsa.Samples.Interrupts .AddElsa() .AddConsoleActivities() .AddHttpActivities() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddActivity() .StartWorkflow(); } diff --git a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Elsa.Samples.AzureServiceBusWorker.csproj b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Elsa.Samples.AzureServiceBusWorker.csproj index cdf4708ca..123854d03 100644 --- a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Elsa.Samples.AzureServiceBusWorker.csproj +++ b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Elsa.Samples.AzureServiceBusWorker.csproj @@ -11,6 +11,8 @@ + + diff --git a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Program.cs b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Program.cs index a83310935..e57ea1d0b 100644 --- a/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Program.cs +++ b/src/samples/worker/Elsa.Samples.AzureServiceBusWorker/Program.cs @@ -1,5 +1,5 @@ using Elsa.Activities.AzureServiceBus.Extensions; -using Elsa.Persistence.InMemory; +using Elsa.Activities.Timers; using Elsa.Persistence.YesSql.Extensions; using Elsa.Samples.AzureServiceBusWorker.Workflows; using Microsoft.Extensions.Configuration; @@ -22,7 +22,7 @@ namespace Elsa.Samples.AzureServiceBusWorker services .AddElsa(options => options.UseYesSqlPersistence()) .AddConsoleActivities() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddAzureServiceBusActivities(options => options.ConnectionString = hostContext.Configuration.GetConnectionString("AzureServiceBus")) .AddWorkflow() .AddWorkflow(); diff --git a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Elsa.Samples.CustomAttributesChildWorker.csproj b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Elsa.Samples.CustomAttributesChildWorker.csproj index df4539f84..7c3080a36 100644 --- a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Elsa.Samples.CustomAttributesChildWorker.csproj +++ b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Elsa.Samples.CustomAttributesChildWorker.csproj @@ -7,11 +7,12 @@ - + - - + + + diff --git a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Program.cs b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Program.cs index 122e03f42..cb744904b 100644 --- a/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Program.cs +++ b/src/samples/worker/Elsa.Samples.CustomAttributesChildWorker/Program.cs @@ -2,6 +2,7 @@ using Elsa.Samples.CustomAttributesChildWorker.Messages; using Elsa.Samples.CustomAttributesChildWorker.Workflows; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Elsa.Activities.Timers; namespace Elsa.Samples.CustomAttributesChildWorker { @@ -19,7 +20,7 @@ namespace Elsa.Samples.CustomAttributesChildWorker { services .AddElsa() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddConsoleActivities() .AddRebusActivities() .AddWorkflow() diff --git a/src/samples/worker/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj b/src/samples/worker/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj index d2b78b64d..d4f04e2ec 100644 --- a/src/samples/worker/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj +++ b/src/samples/worker/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj @@ -7,10 +7,16 @@ - + + + + + + + diff --git a/src/samples/worker/Elsa.Samples.DistributedLock/Program.cs b/src/samples/worker/Elsa.Samples.DistributedLock/Program.cs index 692c449ae..485508a0a 100644 --- a/src/samples/worker/Elsa.Samples.DistributedLock/Program.cs +++ b/src/samples/worker/Elsa.Samples.DistributedLock/Program.cs @@ -1,6 +1,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using System.Threading.Tasks; +using Elsa.Activities.Timers; namespace Elsa.Samples.DistributedLock { @@ -20,7 +21,7 @@ namespace Elsa.Samples.DistributedLock services .AddElsa(options => options.UseRedisLockProvider("localhost:6379,abortConnect=false")) .AddConsoleActivities() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddWorkflow(); }); } diff --git a/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Elsa.Samples.MultiTenantChildWorker.csproj b/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Elsa.Samples.MultiTenantChildWorker.csproj index 40cc8357e..3b0065b81 100644 --- a/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Elsa.Samples.MultiTenantChildWorker.csproj +++ b/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Elsa.Samples.MultiTenantChildWorker.csproj @@ -6,11 +6,12 @@ - + - - + + + diff --git a/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Program.cs b/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Program.cs index d51420ea7..2f4a71e1b 100644 --- a/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Program.cs +++ b/src/samples/worker/Elsa.Samples.MultiTenantChildWorker/Program.cs @@ -2,6 +2,7 @@ using Elsa.Samples.MultiTenantChildWorker.Messages; using Elsa.Samples.MultiTenantChildWorker.Workflows; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Elsa.Activities.Timers; namespace Elsa.Samples.MultiTenantChildWorker { @@ -19,7 +20,7 @@ namespace Elsa.Samples.MultiTenantChildWorker { services .AddElsa() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddConsoleActivities() .AddRebusActivities() .AddWorkflow() diff --git a/src/samples/worker/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj b/src/samples/worker/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj index 0c193f86c..b8b61ae1b 100644 --- a/src/samples/worker/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj +++ b/src/samples/worker/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj @@ -7,11 +7,13 @@ + + diff --git a/src/samples/worker/Elsa.Samples.RebusWorker/Program.cs b/src/samples/worker/Elsa.Samples.RebusWorker/Program.cs index f1346dfca..29abcdb56 100644 --- a/src/samples/worker/Elsa.Samples.RebusWorker/Program.cs +++ b/src/samples/worker/Elsa.Samples.RebusWorker/Program.cs @@ -2,6 +2,7 @@ using Elsa.Samples.RebusWorker.Messages; using Elsa.Samples.RebusWorker.Workflows; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Elsa.Activities.Timers; namespace Elsa.Samples.RebusWorker { @@ -16,7 +17,7 @@ namespace Elsa.Samples.RebusWorker services .AddElsa() .AddConsoleActivities() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddRebusActivities() .AddWorkflow() .AddWorkflow(); diff --git a/src/samples/worker/Elsa.Samples.Timers.Hangfire/CancelTimerWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers.Hangfire/CancelTimerWorkflow.cs new file mode 100644 index 000000000..5dd332977 --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers.Hangfire/CancelTimerWorkflow.cs @@ -0,0 +1,32 @@ +using System.Collections.Generic; + +using Elsa.Activities.Console; +using Elsa.Activities.ControlFlow; +using Elsa.Activities.Timers; +using Elsa.Builders; +using NodaTime; + +namespace Elsa.Samples.Timers +{ + public class CancelTimerWorkflow : IWorkflow + { + public void Build(IWorkflowBuilder workflow) + { + workflow + .StartAt(SystemClock.Instance.GetCurrentInstant().Plus(Duration.FromSeconds(5))) + .WriteLine("CancelTimerWorkflow is executed") + .Then( + activity => activity.Set(x => x.Branches, new HashSet(new[] { "Branch 1", "Branch 2" })), + fork => + { + fork.When("Branch 1").Timer(Duration.FromSeconds(30)).WithId("timer-1").WriteLine("Should not be executed"); + fork.When("Branch 2").Timer(Duration.FromSeconds(10)) + .CancelTimer("timer-1") + .WriteLine("Timer-1 was canceled") + .Timer(Duration.FromSeconds(40)) + .WriteLine("CancelTimerWorkflow finished") + .Finish(); + }); + } + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers.Hangfire/CronTaskWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers.Hangfire/CronTaskWorkflow.cs new file mode 100644 index 000000000..7a8de0527 --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers.Hangfire/CronTaskWorkflow.cs @@ -0,0 +1,18 @@ +using System; +using Elsa.Activities.Console; +using Elsa.Activities.ControlFlow; +using Elsa.Activities.Timers; +using Elsa.Builders; + +namespace Elsa.Samples.Timers +{ + public class CronTaskWorkflow : IWorkflow + { + public void Build(IWorkflowBuilder workflow) + { + workflow + .Cron("0/30 * * * * *") + .WriteLine(() => $"CRON event at {DateTime.Now}"); + } + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers.Hangfire/Elsa.Samples.Timers.Hangfire.csproj b/src/samples/worker/Elsa.Samples.Timers.Hangfire/Elsa.Samples.Timers.Hangfire.csproj new file mode 100644 index 000000000..e6305ab77 --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers.Hangfire/Elsa.Samples.Timers.Hangfire.csproj @@ -0,0 +1,20 @@ + + + + Exe + net5.0 + false + Elsa.Samples.Timers + + + + + + + + + + + + + diff --git a/src/samples/worker/Elsa.Samples.Timers.Hangfire/OneOffWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers.Hangfire/OneOffWorkflow.cs new file mode 100644 index 000000000..463c6ab17 --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers.Hangfire/OneOffWorkflow.cs @@ -0,0 +1,29 @@ +using Elsa.Activities.Console; +using Elsa.Activities.Timers; +using Elsa.Builders; +using NodaTime; + +namespace Elsa.Samples.Timers +{ + /// + /// A workflow that executes only once in the near future. + /// + public class OneOffWorkflow : IWorkflow + { + private readonly Instant _executeAt; + + public OneOffWorkflow(Instant executeAt) + { + _executeAt = executeAt; + } + + public void Build(IWorkflowBuilder workflow) + { + workflow + .StartAt(_executeAt) + .WriteLine(context => $"Started at {context.GetService().GetCurrentInstant()}. Next event happens 30 seconds from now.") + .StartIn(Duration.FromSeconds(30)) + .WriteLine(context => $"Follow-up occurred at {context.GetService().GetCurrentInstant()}."); + } + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers.Hangfire/Program.cs b/src/samples/worker/Elsa.Samples.Timers.Hangfire/Program.cs new file mode 100644 index 000000000..c481d886a --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers.Hangfire/Program.cs @@ -0,0 +1,31 @@ +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using NodaTime; +using Elsa.Activities.Timers.Hangfire.Extensions; +using Hangfire; + +namespace Elsa.Samples.Timers +{ + internal class Program + { + private static async Task Main() => await CreateHostBuilder().UseConsoleLifetime().Build().RunAsync(); + + private static IHostBuilder CreateHostBuilder() => + Host.CreateDefaultBuilder() + .ConfigureServices( + (context, services) => + { + services + .AddElsa() + .AddConsoleActivities() + .AddTimerActivities(options => + options.UseHangfire(configure => + configure.UseSqlServerStorage("Server=(localdb)\\MSSQLLocalDB;Database=ElsaHangfire;Trusted_Connection=True;MultipleActiveResultSets=true"))) + .AddWorkflow() + .AddWorkflow() + .AddWorkflow() + .AddWorkflow(new OneOffWorkflow(SystemClock.Instance.GetCurrentInstant().Plus(Duration.FromSeconds(5)))); + }); + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers.Hangfire/RecurringTaskWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers.Hangfire/RecurringTaskWorkflow.cs new file mode 100644 index 000000000..2bd3840ce --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers.Hangfire/RecurringTaskWorkflow.cs @@ -0,0 +1,27 @@ +using Elsa.Activities.Console; +using Elsa.Activities.Timers; +using Elsa.Builders; +using NodaTime; + +namespace Elsa.Samples.Timers +{ + public class RecurringTaskWorkflow : IWorkflow + { + private readonly IClock _clock; + + public RecurringTaskWorkflow(IClock clock) + { + _clock = clock; + } + + public void Build(IWorkflowBuilder workflow) + { + workflow + .AsSingleton() + .Timer(Duration.FromSeconds(30)) + .WriteLine(context => $"{context.WorkflowExecutionContext.WorkflowInstance.WorkflowInstanceId} triggered by timer at {_clock.GetCurrentInstant()}.") + .Timer(Duration.FromSeconds(30)) + .WriteLine(context => $"{context.WorkflowExecutionContext.WorkflowInstance.WorkflowInstanceId} resumed by timer at {_clock.GetCurrentInstant()}."); + } + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/CancelTimerWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers/CancelTimerWorkflow.cs new file mode 100644 index 000000000..5dd332977 --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers/CancelTimerWorkflow.cs @@ -0,0 +1,32 @@ +using System.Collections.Generic; + +using Elsa.Activities.Console; +using Elsa.Activities.ControlFlow; +using Elsa.Activities.Timers; +using Elsa.Builders; +using NodaTime; + +namespace Elsa.Samples.Timers +{ + public class CancelTimerWorkflow : IWorkflow + { + public void Build(IWorkflowBuilder workflow) + { + workflow + .StartAt(SystemClock.Instance.GetCurrentInstant().Plus(Duration.FromSeconds(5))) + .WriteLine("CancelTimerWorkflow is executed") + .Then( + activity => activity.Set(x => x.Branches, new HashSet(new[] { "Branch 1", "Branch 2" })), + fork => + { + fork.When("Branch 1").Timer(Duration.FromSeconds(30)).WithId("timer-1").WriteLine("Should not be executed"); + fork.When("Branch 2").Timer(Duration.FromSeconds(10)) + .CancelTimer("timer-1") + .WriteLine("Timer-1 was canceled") + .Timer(Duration.FromSeconds(40)) + .WriteLine("CancelTimerWorkflow finished") + .Finish(); + }); + } + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/Elsa.Samples.Timers.csproj b/src/samples/worker/Elsa.Samples.Timers/Elsa.Samples.Timers.Quartz.csproj similarity index 68% rename from src/samples/worker/Elsa.Samples.Timers/Elsa.Samples.Timers.csproj rename to src/samples/worker/Elsa.Samples.Timers/Elsa.Samples.Timers.Quartz.csproj index bc261bc77..6f18de34e 100644 --- a/src/samples/worker/Elsa.Samples.Timers/Elsa.Samples.Timers.csproj +++ b/src/samples/worker/Elsa.Samples.Timers/Elsa.Samples.Timers.Quartz.csproj @@ -4,6 +4,7 @@ Exe net5.0 false + Elsa.Samples.Timers @@ -11,6 +12,8 @@ + + diff --git a/src/samples/worker/Elsa.Samples.Timers/Program.cs b/src/samples/worker/Elsa.Samples.Timers/Program.cs index d788fcbb6..c799ecba6 100644 --- a/src/samples/worker/Elsa.Samples.Timers/Program.cs +++ b/src/samples/worker/Elsa.Samples.Timers/Program.cs @@ -3,6 +3,7 @@ using Elsa.Persistence.YesSql.Extensions; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using NodaTime; +using Elsa.Activities.Timers; namespace Elsa.Samples.Timers { @@ -18,8 +19,9 @@ namespace Elsa.Samples.Timers services .AddElsa(options => options.UseYesSqlPersistence()) .AddConsoleActivities() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddWorkflow() + .AddWorkflow() //.AddWorkflow() //.AddWorkflow(new OneOffWorkflow(SystemClock.Instance.GetCurrentInstant().Plus(Duration.FromSeconds(5)))) ; diff --git a/src/samples/worker/Elsa.Samples.WhileLoopWorker/Elsa.Samples.WhileLoopWorker.csproj b/src/samples/worker/Elsa.Samples.WhileLoopWorker/Elsa.Samples.WhileLoopWorker.csproj index 5d05f57bb..df6bbf46f 100644 --- a/src/samples/worker/Elsa.Samples.WhileLoopWorker/Elsa.Samples.WhileLoopWorker.csproj +++ b/src/samples/worker/Elsa.Samples.WhileLoopWorker/Elsa.Samples.WhileLoopWorker.csproj @@ -10,6 +10,8 @@ + + diff --git a/src/samples/worker/Elsa.Samples.WhileLoopWorker/Program.cs b/src/samples/worker/Elsa.Samples.WhileLoopWorker/Program.cs index e35cd2afb..2232f1c9d 100644 --- a/src/samples/worker/Elsa.Samples.WhileLoopWorker/Program.cs +++ b/src/samples/worker/Elsa.Samples.WhileLoopWorker/Program.cs @@ -20,7 +20,7 @@ namespace Elsa.Samples.WhileLoopWorker services .AddElsa(options => options.UseYesSqlPersistence()) .AddConsoleActivities() - .AddTimerActivities() + .AddTimerActivities(options => options.UseQuartzProvider()) .AddSingleton() .AddHostedService() .AddActivity()