From b8d313d0c2d4753ad8ef8f33c55aa8f66d61f95d Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 6 May 2021 22:15:05 +0200 Subject: [PATCH] Fix unscheduling of triggers upon publishing/retracting workflows --- .../Handlers/RemoveScheduledTriggers.cs | 5 +- .../Services/IWorkflowScheduler.cs | 1 + .../BackgroundJobClientExtensions.cs | 18 ++++++- .../Services/HangfireWorkflowScheduler.cs | 6 +++ .../Services/QuartzWorkflowScheduler.cs | 18 ++++++- .../DbContextOptionsBuilderExtensions.cs | 2 +- .../Extensions/ServiceCollectionExtensions.cs | 2 +- .../Elsa.Samples.Server.Host.csproj | 1 + .../Elsa.Samples.Server.Host/Startup.cs | 15 ++++-- .../Workflows/ConditionWorkflow.cs | 27 ---------- .../Workflows/FaultyWorkflow.cs | 51 ------------------- .../Workflows/GoodbyeWorld.cs | 17 ------- .../Workflows/HeartbeatWorkflow.cs | 21 -------- .../Workflows/NamingWorkflow.cs | 31 ----------- 14 files changed, 59 insertions(+), 156 deletions(-) delete mode 100644 src/samples/server/Elsa.Samples.Server.Host/Workflows/ConditionWorkflow.cs delete mode 100644 src/samples/server/Elsa.Samples.Server.Host/Workflows/FaultyWorkflow.cs delete mode 100644 src/samples/server/Elsa.Samples.Server.Host/Workflows/GoodbyeWorld.cs delete mode 100644 src/samples/server/Elsa.Samples.Server.Host/Workflows/HeartbeatWorkflow.cs delete mode 100644 src/samples/server/Elsa.Samples.Server.Host/Workflows/NamingWorkflow.cs diff --git a/src/activities/Elsa.Activities.Temporal.Common/Handlers/RemoveScheduledTriggers.cs b/src/activities/Elsa.Activities.Temporal.Common/Handlers/RemoveScheduledTriggers.cs index e931a8486..ef276c4b3 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/Handlers/RemoveScheduledTriggers.cs +++ b/src/activities/Elsa.Activities.Temporal.Common/Handlers/RemoveScheduledTriggers.cs @@ -7,7 +7,7 @@ using MediatR; namespace Elsa.Activities.Temporal.Common.Handlers { - public class RemoveScheduledTriggers : INotificationHandler + public class RemoveScheduledTriggers : INotificationHandler, INotificationHandler, INotificationHandler { private readonly string[] _supportedTypes = { nameof(Timer), nameof(StartAt), nameof(Cron) }; private readonly IWorkflowScheduler _workflowScheduler; @@ -25,5 +25,8 @@ namespace Elsa.Activities.Temporal.Common.Handlers notification.WorkflowExecutionContext.WorkflowInstance.TenantId, cancellationToken); } + + public Task Handle(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) => _workflowScheduler.UnscheduleWorkflowDefinitionAsync(notification.WorkflowDefinition.DefinitionId, notification.WorkflowDefinition.TenantId, cancellationToken); + public Task Handle(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => _workflowScheduler.UnscheduleWorkflowDefinitionAsync(notification.WorkflowDefinition.DefinitionId, notification.WorkflowDefinition.TenantId, cancellationToken); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Temporal.Common/Services/IWorkflowScheduler.cs b/src/activities/Elsa.Activities.Temporal.Common/Services/IWorkflowScheduler.cs index 2994e8f22..4ec6eb131 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/Services/IWorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Common/Services/IWorkflowScheduler.cs @@ -9,5 +9,6 @@ namespace Elsa.Activities.Temporal.Common.Services Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, Instant startAt, Duration? interval = default, CancellationToken cancellationToken = default); Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string cronExpression, CancellationToken cancellationToken = default); Task UnscheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, CancellationToken cancellationToken = default); + Task UnscheduleWorkflowDefinitionAsync(string workflowDefinitionId, string? tenantId, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/BackgroundJobClientExtensions.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/BackgroundJobClientExtensions.cs index 69bd313d6..3cee16cc1 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/BackgroundJobClientExtensions.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Extensions/BackgroundJobClientExtensions.cs @@ -25,7 +25,23 @@ namespace Elsa.Activities.Temporal .Where(x => x.Value.Job.Type == workflowJobType && ((RunHangfireWorkflowJobModel)x.Value.Job.Args[0]).GetIdentity() == identity); foreach (var job in jobs) - BackgroundJob.Delete(job.Key); + 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.Hangfire/Services/HangfireWorkflowScheduler.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs index 2fdb0a360..6c425bf9a 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs @@ -43,6 +43,12 @@ namespace Elsa.Activities.Temporal.Hangfire.Services return Task.CompletedTask; } + public Task UnscheduleWorkflowDefinitionAsync(string workflowDefinitionId, string? tenantId, CancellationToken cancellationToken = default) + { + _backgroundJobClient.UnscheduleJobWhenAlreadyExists(workflowDefinitionId, tenantId); + return Task.CompletedTask; + } + private RunHangfireWorkflowJobModel CreateData(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string? cronExpression = null) { return new( diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowScheduler.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowScheduler.cs index 62fcebd0a..d7be49241 100644 --- a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowScheduler.cs @@ -5,6 +5,8 @@ using Elsa.Activities.Temporal.Quartz.Jobs; using Microsoft.Extensions.Logging; using NodaTime; using Quartz; +using Quartz.Impl.Matchers; +using Quartz.Util; namespace Elsa.Activities.Temporal.Quartz.Services { @@ -47,6 +49,16 @@ namespace Elsa.Activities.Temporal.Quartz.Services if (existingTrigger != null) await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken); } + + public async Task UnscheduleWorkflowDefinitionAsync(string workflowDefinitionId, string? tenantId, CancellationToken cancellationToken) + { + var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var groupName = CreateTriggerGroupKey(tenantId, workflowDefinitionId); + var existingTriggers = await scheduler.GetTriggerKeys(GroupMatcher.GroupEquals(groupName), cancellationToken); + + foreach (var existingTrigger in existingTriggers) + await scheduler.UnscheduleJob(existingTrigger, cancellationToken); + } private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) { @@ -59,7 +71,7 @@ namespace Elsa.Activities.Temporal.Quartz.Services if (existingTrigger != null) await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken); - + await scheduler.ScheduleJob(trigger, cancellationToken); } catch (SchedulerException e) @@ -83,8 +95,10 @@ namespace Elsa.Activities.Temporal.Quartz.Services private TriggerKey CreateTriggerKey(string? tenantId, string? workflowDefinitionId, string? workflowInstanceId, string activityId) { - var groupName = $"tenant:{tenantId ?? "default"}-workflow-instance:{workflowInstanceId ?? workflowDefinitionId}"; + var groupName = CreateTriggerGroupKey(tenantId, workflowInstanceId ?? workflowDefinitionId); return new TriggerKey($"activity:{activityId}", groupName); } + + private string CreateTriggerGroupKey(string? tenantId, string? workflowId) => $"tenant:{tenantId ?? "default"}-workflow:{workflowId}"; } } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Sqlite/DbContextOptionsBuilderExtensions.cs b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Sqlite/DbContextOptionsBuilderExtensions.cs index 6f5218e9b..54b310741 100644 --- a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Sqlite/DbContextOptionsBuilderExtensions.cs +++ b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Sqlite/DbContextOptionsBuilderExtensions.cs @@ -4,6 +4,6 @@ namespace Elsa.Persistence.EntityFramework.Sqlite { public static class DbContextOptionsBuilderExtensions { - public static DbContextOptionsBuilder UseSqlite(this DbContextOptionsBuilder builder) => builder.UseSqlite("Data Source=elsa.db;Cache=Shared;", db => db.MigrationsAssembly(typeof(SqliteElsaContextFactory).Assembly.GetName().Name)); + public static DbContextOptionsBuilder UseSqlite(this DbContextOptionsBuilder builder) => builder.UseSqlite("Data Source=elsa.sqlite.db;Cache=Shared;", db => db.MigrationsAssembly(typeof(SqliteElsaContextFactory).Assembly.GetName().Name)); } } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.YesSql/Extensions/ServiceCollectionExtensions.cs b/src/persistence/Elsa.Persistence.YesSql/Extensions/ServiceCollectionExtensions.cs index bdcedd45d..a35b078e5 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Extensions/ServiceCollectionExtensions.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Extensions/ServiceCollectionExtensions.cs @@ -16,7 +16,7 @@ namespace Elsa.Persistence.YesSql { public static class ServiceCollectionExtensions { - public static ElsaOptionsBuilder UseYesSqlPersistence(this ElsaOptionsBuilder elsa) => elsa.UseYesSqlPersistence(config => config.UseSqLite("Data Source=elsa.db;Cache=Shared", IsolationLevel.ReadUncommitted)); + public static ElsaOptionsBuilder UseYesSqlPersistence(this ElsaOptionsBuilder elsa) => elsa.UseYesSqlPersistence(config => config.UseSqLite("Data Source=elsa.yessql.db;Cache=Shared", IsolationLevel.ReadUncommitted)); public static ElsaOptionsBuilder UseYesSqlPersistence(this ElsaOptionsBuilder elsa, Action configure) => elsa.UseYesSqlPersistence((_, config) => configure(config)); public static ElsaOptionsBuilder UseYesSqlPersistence(this ElsaOptionsBuilder elsa, Action configure) 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 3ab387071..4b5f1b101 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 @@ -28,6 +28,7 @@ + diff --git a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs index 98da35742..86a1142b0 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs +++ b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs @@ -2,6 +2,7 @@ using Elsa.Activities.UserTask.Extensions; using Elsa.Persistence.EntityFramework.Core.Extensions; using Elsa.Persistence.EntityFramework.Sqlite; using Elsa.Persistence.EntityFramework.SqlServer; +using Elsa.Persistence.YesSql; using Elsa.Samples.Server.Host.Activities; using Elsa.Server.Hangfire.Extensions; using Hangfire; @@ -12,6 +13,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using NodaTime; using NodaTime.Serialization.JsonNet; +using Quartz; namespace Elsa.Samples.Server.Host { @@ -52,8 +54,8 @@ namespace Elsa.Samples.Server.Host .AddActivityPropertyOptionsProvider() .AddRuntimeSelectItemsProvider() .AddElsa(elsa => elsa - //.UseNonPooledEntityFrameworkPersistence(ef => ef.UseSqlite()) - .UseEntityFrameworkPersistence(ef => ef.UseSqlServer("Server=LAPTOP-B76STK67;Database=Elsa;Integrated Security=true;MultipleActiveResultSets=True;")) + .UseEntityFrameworkPersistence(ef => ef.UseSqlite()) + //.UseEntityFrameworkPersistence(ef => ef.UseSqlServer("Server=LAPTOP-B76STK67;Database=Elsa;Integrated Security=true;MultipleActiveResultSets=True;")) //.UseYesSqlPersistence() // Using Hangfire as the dispatcher for workflow execution in the background. @@ -63,10 +65,17 @@ namespace Elsa.Samples.Server.Host .AddHttpActivities(elsaSection.GetSection("Http").Bind) .AddEmailActivities(elsaSection.GetSection("Smtp").Bind) .AddQuartzTemporalActivities() + // .AddQuartzTemporalActivities(configureQuartz: quartz => quartz.UsePersistentStore(x => + // { + // x.UseJsonSerializer(); + // x.UseGenericDatabase("SqlServer", ado => + // { + // ado.ConnectionString = "Server=LAPTOP-B76STK67;Database=Elsa;Integrated Security=true;MultipleActiveResultSets=True;"; + // }); + // })) .AddJavaScriptActivities() .AddUserTaskActivities() .AddActivitiesFrom() - .AddWorkflowsFrom() ); // Elsa API endpoints. diff --git a/src/samples/server/Elsa.Samples.Server.Host/Workflows/ConditionWorkflow.cs b/src/samples/server/Elsa.Samples.Server.Host/Workflows/ConditionWorkflow.cs deleted file mode 100644 index f0e60a0e0..000000000 --- a/src/samples/server/Elsa.Samples.Server.Host/Workflows/ConditionWorkflow.cs +++ /dev/null @@ -1,27 +0,0 @@ -using System; -using Elsa.Activities.Console; -using Elsa.Activities.ControlFlow; -using Elsa.Activities.Primitives; -using Elsa.Activities.Temporal; -using Elsa.Builders; -using NodaTime; - -namespace Elsa.Samples.Server.Host.Workflows -{ - public class ConditionWorkflow : IWorkflow - { - public void Build(IWorkflowBuilder builder) - { - builder - .WithDisplayName("Conditions") - .Then(() => Console.WriteLine("What is your age?")).WithDisplayName("Write").WithDescription("What is your age?") - .ReadLine() - .Timer(Duration.FromMinutes(5)) - .SetVariable("Age", context => int.Parse(context.GetInput()!)) - .If( - context => context.GetVariable("Age") < 18, - whenTrue => whenTrue.WriteLine("You are not allowed to drink beer."), - whenFalse => whenFalse.WriteLine("Enjoy your beer!")); - } - } -} \ No newline at end of file diff --git a/src/samples/server/Elsa.Samples.Server.Host/Workflows/FaultyWorkflow.cs b/src/samples/server/Elsa.Samples.Server.Host/Workflows/FaultyWorkflow.cs deleted file mode 100644 index 22b3c9b08..000000000 --- a/src/samples/server/Elsa.Samples.Server.Host/Workflows/FaultyWorkflow.cs +++ /dev/null @@ -1,51 +0,0 @@ -using System; -using System.Net; -using Elsa.Activities.Http; -using Elsa.Activities.Http.Models; -using Elsa.Builders; -using Elsa.Serialization; -using Elsa.Services.Models; - -namespace Elsa.Samples.Server.Host.Workflows -{ - /// - /// A workflow that is triggered when HTTP requests are made to /hello and writes a response. - /// - public class FaultyWorkflow : IWorkflow - { - private readonly IContentSerializer _contentSerializer; - - public FaultyWorkflow(IContentSerializer contentSerializer) - { - _contentSerializer = contentSerializer; - } - - public void Build(IWorkflowBuilder builder) - { - builder - .HttpEndpoint("/faulty") - .Then(MaybeThrow) - .WriteHttpResponse(response => response.WithStatusCode(HttpStatusCode.OK).WithContentType("application/json").WithContent(WriteWorkflowInfoAsync)); - } - - private void MaybeThrow(ActivityExecutionContext context) - { - var model = context.GetInput()!; - var fault = model.QueryString.GetItem("fault") == "true"; - - if (fault) - throw new Exception("This is quite a serious fault!"); - } - - private string WriteWorkflowInfoAsync(ActivityExecutionContext context) - { - var model = new - { - WorkflowInstanceId = context.WorkflowInstance.Id, - WorkflowStatus = context.WorkflowInstance.WorkflowStatus - }; - - return _contentSerializer.Serialize(model); - } - } -} \ No newline at end of file diff --git a/src/samples/server/Elsa.Samples.Server.Host/Workflows/GoodbyeWorld.cs b/src/samples/server/Elsa.Samples.Server.Host/Workflows/GoodbyeWorld.cs deleted file mode 100644 index 04218e394..000000000 --- a/src/samples/server/Elsa.Samples.Server.Host/Workflows/GoodbyeWorld.cs +++ /dev/null @@ -1,17 +0,0 @@ -using System.Net; -using Elsa.Activities.Http; -using Elsa.Builders; - -namespace Elsa.Samples.Server.Host.Workflows -{ - public class GoodbyeWorld : IWorkflow - { - public void Build(IWorkflowBuilder builder) - { - builder - .WithDisplayName("Goodbye cruel World!") - .HttpEndpoint("/goodbye-world") - .WriteHttpResponse(HttpStatusCode.OK, "Goodbye Cruel World!", "text/plain"); - } - } -} \ No newline at end of file diff --git a/src/samples/server/Elsa.Samples.Server.Host/Workflows/HeartbeatWorkflow.cs b/src/samples/server/Elsa.Samples.Server.Host/Workflows/HeartbeatWorkflow.cs deleted file mode 100644 index 5a343e4c0..000000000 --- a/src/samples/server/Elsa.Samples.Server.Host/Workflows/HeartbeatWorkflow.cs +++ /dev/null @@ -1,21 +0,0 @@ -using Elsa.Activities.Console; -using Elsa.Activities.Temporal; -using Elsa.Builders; -using NodaTime; - -namespace Elsa.Samples.Server.Host.Workflows -{ - public class HeartbeatWorkflow : IWorkflow - { - private readonly IClock _clock; - public HeartbeatWorkflow(IClock clock) => _clock = clock; - - public void Build(IWorkflowBuilder builder) - { - builder - .WithDisplayName("Timer") - .Timer(Duration.FromSeconds(10)) - .WriteLine(() => $"Heartbeat at {_clock.GetCurrentInstant()}"); - } - } -} \ No newline at end of file diff --git a/src/samples/server/Elsa.Samples.Server.Host/Workflows/NamingWorkflow.cs b/src/samples/server/Elsa.Samples.Server.Host/Workflows/NamingWorkflow.cs deleted file mode 100644 index 822530302..000000000 --- a/src/samples/server/Elsa.Samples.Server.Host/Workflows/NamingWorkflow.cs +++ /dev/null @@ -1,31 +0,0 @@ -using System; -using System.Net; -using Elsa.Activities.ControlFlow; -using Elsa.Activities.Http; -using Elsa.Activities.Http.Models; -using Elsa.Activities.Primitives; -using Elsa.Builders; -using Microsoft.AspNetCore.Server.Kestrel.Core.Internal.Http; - -namespace Elsa.Samples.Server.Host.Workflows -{ - public class NamingWorkflow : IWorkflow - { - public void Build(IWorkflowBuilder builder) - { - builder - .WithDisplayName("Onboarding") - .HttpEndpoint("/signup") - .Correlate(() => Guid.NewGuid().ToString("N")) - .WriteHttpResponse(x => x.WithStatusCode(HttpStatusCode.OK).WithContent(context => $"Tell me your name please. Use correlation ID {context.WorkflowExecutionContext.CorrelationId}").WithContentType("text/plain")) - .HttpEndpoint(x => x.WithPath("/signup").WithMethod(HttpMethod.Post.ToString()).WithReadContent()) - .SetVariable("Name", context => (string) context.GetInput()!.Body!) - .SetName(context => context.GetVariable("Name")) - - .WriteHttpResponse(x => x - .WithStatusCode(HttpStatusCode.OK) - .WithContent(context => $"Welcome aboard, {context.GetVariable("Name")}!") - .WithContentType("text/plain")); - } - } -} \ No newline at end of file