Fix unscheduling of triggers upon publishing/retracting workflows

This commit is contained in:
Sipke Schoorstra 2021-05-06 22:15:05 +02:00
parent cab4a24713
commit b8d313d0c2
14 changed files with 59 additions and 156 deletions

View file

@ -7,7 +7,7 @@ using MediatR;
namespace Elsa.Activities.Temporal.Common.Handlers
{
public class RemoveScheduledTriggers : INotificationHandler<BlockingActivityRemoved>
public class RemoveScheduledTriggers : INotificationHandler<BlockingActivityRemoved>, INotificationHandler<WorkflowDefinitionPublished>, INotificationHandler<WorkflowDefinitionRetracted>
{
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);
}
}

View file

@ -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);
}
}

View file

@ -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);
}
}
}

View file

@ -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(

View file

@ -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<TriggerKey>.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}";
}
}

View file

@ -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));
}
}

View file

@ -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<IConfiguration> configure) => elsa.UseYesSqlPersistence((_, config) => configure(config));
public static ElsaOptionsBuilder UseYesSqlPersistence(this ElsaOptionsBuilder elsa, Action<IServiceProvider, IConfiguration> configure)

View file

@ -28,6 +28,7 @@
<ItemGroup>
<PackageReference Include="Hangfire.InMemory" Version="0.3.4" />
<PackageReference Include="Quartz.Serialization.Json" Version="3.3.2" />
</ItemGroup>
</Project>

View file

@ -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<VehicleActivity>()
.AddRuntimeSelectItemsProvider<VehicleActivity>()
.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<VehicleActivity>()
.AddWorkflowsFrom<Startup>()
);
// Elsa API endpoints.

View file

@ -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<string>()!))
.If(
context => context.GetVariable<int>("Age") < 18,
whenTrue => whenTrue.WriteLine("You are not allowed to drink beer."),
whenFalse => whenFalse.WriteLine("Enjoy your beer!"));
}
}
}

View file

@ -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
{
/// <summary>
/// A workflow that is triggered when HTTP requests are made to /hello and writes a response.
/// </summary>
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<HttpRequestModel>()!;
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);
}
}
}

View file

@ -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");
}
}
}

View file

@ -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()}");
}
}
}

View file

@ -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<HttpRequestModel>()!.Body!)
.SetName(context => context.GetVariable<string>("Name"))
.WriteHttpResponse(x => x
.WithStatusCode(HttpStatusCode.OK)
.WithContent(context => $"Welcome aboard, {context.GetVariable<string>("Name")}!")
.WithContentType("text/plain"));
}
}
}