diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/BackgroundJobClientExtensions.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/BackgroundJobClientExtensions.cs deleted file mode 100644 index 3cee16cc1..000000000 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/BackgroundJobClientExtensions.cs +++ /dev/null @@ -1,47 +0,0 @@ -using System; -using System.Linq; -using Elsa.Activities.Temporal.Hangfire.Jobs; -using Elsa.Activities.Temporal.Hangfire.Models; -using Hangfire; - -// ReSharper disable once CheckNamespace -namespace Elsa.Activities.Temporal -{ - public static class BackgroundJobClientExtensions - { - 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) - backgroundJobClient.Delete(job.Key); - } - - public static void UnscheduleJobWhenAlreadyExists(this IBackgroundJobClient backgroundJobClient, string workflowDefinitionId, string? tenantId) - { - var monitor = JobStorage.Current.GetMonitoringApi(); - var workflowJobType = typeof(RunHangfireWorkflowJob); - - var jobs = monitor.ScheduledJobs(0, int.MaxValue) - .Where(x => - { - var model = (RunHangfireWorkflowJobModel) x.Value.Job.Args[0]; - return x.Value.Job.Type == workflowJobType && model.WorkflowDefinitionId == workflowDefinitionId && model.TenantId == tenantId; - }); - - foreach (var job in jobs) - backgroundJobClient.Delete(job.Key); - } - } -} diff --git a/src/activities/Elsa.Activities.Temporal.Common/Extensions/DurationExtensions.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/DurationExtensions.cs similarity index 67% rename from src/activities/Elsa.Activities.Temporal.Common/Extensions/DurationExtensions.cs rename to src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/DurationExtensions.cs index de709a674..20be8f4a7 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/Extensions/DurationExtensions.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/DurationExtensions.cs @@ -1,14 +1,12 @@ -// ReSharper disable once CheckNamespace -namespace NodaTime +using NodaTime; + +namespace Elsa.Activities.Temporal.Hangfire.Extensions { public static class DurationExtensions { public static string ToCronExpression(this Duration duration) { - static string CreateCronComponent(int number) - { - return (number > 0 ? $"*/{number}" : "* "); - } + static string CreateCronComponent(int number) => (number > 0 ? $"*/{number}" : "*"); var cron = CreateCronComponent(duration.Seconds); cron += ' ' + CreateCronComponent(duration.Minutes); diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/ElsaOptionsExtensions.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/ElsaOptionsExtensions.cs index 704d77407..c7d61c85a 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/ElsaOptionsExtensions.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/ElsaOptionsExtensions.cs @@ -13,13 +13,15 @@ namespace Elsa /// Elsa options /// A Hangfire configuration callback /// The Elsa options, enabling method chaining + /// Configure Hangfire job server settings public static ElsaOptionsBuilder AddHangfireTemporalActivities( this ElsaOptionsBuilder options, - Action configure) => - options.AddCommonTemporalActivities(timer => timer.UseHangfire(configure)); + Action configure, + Action? configureJobServer = default) => + options.AddCommonTemporalActivities(timer => timer.UseHangfire(configure, configureJobServer)); /// - /// Adds temporal (time-based) activities to Elsa without using the Hangfire implementation. You need to add Hangfire services yourself + /// Adds temporal (time-based) activities to Elsa without using the Hangfire implementation. You need to add Hangfire services yourself. /// /// Elsa options /// The Elsa options, enabling method chaining diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/TimersOptionsExtensions.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/TimersOptionsExtensions.cs index 21ba35bf9..e6a67192e 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/TimersOptionsExtensions.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/TimersOptionsExtensions.cs @@ -1,4 +1,5 @@ using System; +using Elsa.Activities.Temporal; using Elsa.Activities.Temporal.Common.Options; using Elsa.Activities.Temporal.Common.Services; using Elsa.Activities.Temporal.Hangfire.Services; @@ -19,7 +20,8 @@ namespace Elsa { timersOptions.Services .AddSingleton() - .AddSingleton(); + .AddSingleton() + .AddSingleton(); } /// @@ -30,15 +32,21 @@ namespace Elsa /// /// The TimersOptions being configured /// Configure Hangfire settings - public static void UseHangfire(this TimersOptions timersOptions, Action configure) + /// Configure Hangfire job server settings + public static void UseHangfire(this TimersOptions timersOptions, Action configure, Action? configureJobServer = default) { timersOptions.UseHangfire(); + var services = timersOptions.Services; + // Add Hangfire services. - timersOptions.Services.AddHangfire(configure); + services.AddHangfire(configure); // Add the processing server as IHostedService - timersOptions.Services.AddHangfireServer(); + if(configureJobServer != null) + services.AddHangfireServer(configureJobServer); + else + services.AddHangfireServer(); } } } diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs index 6c425bf9a..27509c53d 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs @@ -1,6 +1,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Temporal.Common.Services; +using Elsa.Activities.Temporal.Hangfire.Extensions; using Elsa.Activities.Temporal.Hangfire.Models; using Hangfire; using NodaTime; @@ -9,43 +10,46 @@ namespace Elsa.Activities.Temporal.Hangfire.Services { public class HangfireWorkflowScheduler : IWorkflowScheduler { - private readonly IBackgroundJobClient _backgroundJobClient; - private readonly ICrontabParser _crontabParser; + private readonly JobManager _jobManager; - public HangfireWorkflowScheduler(IBackgroundJobClient backgroundJobClient, ICrontabParser crontabParser) + public HangfireWorkflowScheduler(JobManager jobManager) { - _backgroundJobClient = backgroundJobClient; - _crontabParser = crontabParser; + _jobManager = jobManager; } public Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, Instant startAt, Duration? interval, CancellationToken cancellationToken) { - var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, interval?.ToCronExpression()); + var cron = interval?.ToCronExpression(); + var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, cron); - _backgroundJobClient.ScheduleWorkflow(data, startAt.ToDateTimeOffset()); + _jobManager.ScheduleJob(data, startAt); + + if (cron != null) + _jobManager.ScheduleJob(data, cron); return Task.CompletedTask; } - + public Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string cronExpression, CancellationToken cancellationToken) { var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, cronExpression); - var instant = _crontabParser.GetNextOccurrence(cronExpression); - _backgroundJobClient.ScheduleWorkflow(data, instant.ToDateTimeOffset()); + _jobManager.ScheduleJob(data, cronExpression); return Task.CompletedTask; } public Task UnscheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, CancellationToken cancellationToken = default) { - _backgroundJobClient.UnscheduleJobWhenAlreadyExists(CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId)); + var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId); + var identity = data.GetIdentity(); + _jobManager.UnscheduleJob(identity); return Task.CompletedTask; } public Task UnscheduleWorkflowDefinitionAsync(string workflowDefinitionId, string? tenantId, CancellationToken cancellationToken = default) { - _backgroundJobClient.UnscheduleJobWhenAlreadyExists(workflowDefinitionId, tenantId); + _jobManager.UnscheduleJobs(workflowDefinitionId, tenantId); return Task.CompletedTask; } diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/JobManager.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/JobManager.cs new file mode 100644 index 000000000..62062010a --- /dev/null +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/JobManager.cs @@ -0,0 +1,71 @@ +using System.Collections.Generic; +using System.Linq; +using Elsa.Activities.Temporal.Hangfire.Jobs; +using Elsa.Activities.Temporal.Hangfire.Models; +using Hangfire; +using Hangfire.Common; +using Hangfire.Storage; +using NodaTime; + +namespace Elsa.Activities.Temporal.Hangfire.Services +{ + public class JobManager + { + private readonly IBackgroundJobClient _backgroundJobClient; + private readonly IRecurringJobManager _recurringJobManager; + private readonly JobStorage _jobStorage; + + public JobManager(IBackgroundJobClient backgroundJobClient, IRecurringJobManager recurringJobManager, JobStorage jobStorage) + { + _backgroundJobClient = backgroundJobClient; + _recurringJobManager = recurringJobManager; + _jobStorage = jobStorage; + } + + public void ScheduleJob(RunHangfireWorkflowJobModel data, Instant instant) + { + UnscheduleJob(data); + _backgroundJobClient.Schedule(job => job.ExecuteAsync(data), instant.ToDateTimeOffset()); + } + + public void ScheduleJob(RunHangfireWorkflowJobModel data, string cronExpression) + { + var identity = data.GetIdentity(); + _recurringJobManager.AddOrUpdate(identity, job => job.ExecuteAsync(data), cronExpression); + } + + public void UnscheduleJob(RunHangfireWorkflowJobModel data) + { + var identity = data.GetIdentity(); + UnscheduleJob(identity); + } + + public void UnscheduleJob(string identity) + { + var job = QueryJobs().FirstOrDefault(x => GetJobModel(x.Job).GetIdentity() == identity); + + if (job == null) + return; + + _backgroundJobClient.Delete(job.Id); + _recurringJobManager.RemoveIfExists(job.Id); + } + + public void UnscheduleJobs(string workflowDefinitionId, string? tenantId) + { + var jobs = QueryJobs(); + + var jobsToRemove = jobs.Where(x => + { + var model = GetJobModel(x.Job); + return model.WorkflowDefinitionId == workflowDefinitionId && model.TenantId == tenantId; + }).ToList(); + + foreach (var job in jobsToRemove) + _recurringJobManager.RemoveIfExists(job.Id); + } + + private RunHangfireWorkflowJobModel GetJobModel(Job job) => (RunHangfireWorkflowJobModel) job.Args[0]; + private IEnumerable QueryJobs() => _jobStorage.GetConnection().GetRecurringJobs().Where(x => x.Job.Type == typeof(RunHangfireWorkflowJob)); + } +} \ No newline at end of file diff --git a/src/designer/elsa-workflows-studio/src/components.d.ts b/src/designer/elsa-workflows-studio/src/components.d.ts index e7b95eb19..94cd400cf 100644 --- a/src/designer/elsa-workflows-studio/src/components.d.ts +++ b/src/designer/elsa-workflows-studio/src/components.d.ts @@ -37,7 +37,7 @@ export namespace Components { } interface ElsaDesignerTree { "activityContextMenu"?: ActivityContextMenuState; - "activityContextMenuButton"?: string; + "activityContextMenuButton"?: (activity: ActivityModel) => string; "mode": WorkflowDesignerMode; "model": WorkflowModel; "removeActivity": (activity: ActivityModel) => Promise; @@ -544,7 +544,7 @@ declare namespace LocalJSX { } interface ElsaDesignerTree { "activityContextMenu"?: ActivityContextMenuState; - "activityContextMenuButton"?: string; + "activityContextMenuButton"?: (activity: ActivityModel) => string; "mode"?: WorkflowDesignerMode; "model"?: WorkflowModel; "onActivityContextMenuButtonClicked"?: (event: CustomEvent) => void; diff --git a/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj b/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj index 2570b8418..091ecd4f3 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj +++ b/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj @@ -10,6 +10,7 @@ + diff --git a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs index 678e2b5b1..ed9d9f0cc 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs +++ b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs @@ -1,9 +1,11 @@ +using System; using Elsa.Activities.UserTask.Extensions; using Elsa.Persistence.EntityFramework.Core.Extensions; using Elsa.Persistence.EntityFramework.PostgreSql; using Elsa.Persistence.EntityFramework.Sqlite; using Elsa.Persistence.YesSql; using Elsa.Samples.Server.Host.Activities; +using Hangfire; using Microsoft.AspNetCore.Builder; using Microsoft.AspNetCore.Hosting; using Microsoft.AspNetCore.Mvc.Versioning; @@ -38,7 +40,8 @@ namespace Elsa.Samples.Server.Host .AddConsoleActivities() .AddHttpActivities(elsaSection.GetSection("Http").Bind) .AddEmailActivities(elsaSection.GetSection("Smtp").Bind) - .AddQuartzTemporalActivities() + //.AddQuartzTemporalActivities() + .AddHangfireTemporalActivities(hangfire => hangfire.UseInMemoryStorage(), (_, hangfireServer) => hangfireServer.SchedulePollingInterval = TimeSpan.FromSeconds(5)) .AddJavaScriptActivities() .AddUserTaskActivities() .AddActivitiesFrom() diff --git a/src/samples/server/Elsa.Samples.Server.Host/appsettings.Development.json b/src/samples/server/Elsa.Samples.Server.Host/appsettings.Development.json index 9c8e47589..751217e22 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/appsettings.Development.json +++ b/src/samples/server/Elsa.Samples.Server.Host/appsettings.Development.json @@ -1,9 +1,9 @@ { "Logging": { "LogLevel": { - "Default": "Debug", + "Default": "Warning", "Orleans": "Warning", - "System": "Information", + "System": "Warning", "Microsoft": "Warning" } }