Fix workflow deletion and updates handling

This commit is contained in:
Raymond den Haan 2024-05-17 10:34:43 +02:00 committed by raymonddenhaan
parent 0a86729040
commit 5dcf6f382a
12 changed files with 221 additions and 9 deletions

View file

@ -1,4 +1,3 @@
using Elsa.Common.DistributedLocks.Noop;
using Elsa.EntityFrameworkCore.Extensions;
using Elsa.EntityFrameworkCore.Modules.Management;
using Elsa.EntityFrameworkCore.Modules.Runtime;
@ -7,6 +6,7 @@ using Elsa.Extensions;
using Elsa.ServerAndStudio.Web.Extensions;
using Elsa.MassTransit.Extensions;
using Elsa.ServerAndStudio.Web.Enums;
using Medallion.Threading.FileSystem;
using Microsoft.AspNetCore.Mvc;
using Microsoft.Data.Sqlite;
using Proto.Persistence.Sqlite;
@ -73,7 +73,7 @@ services
});
}
runtime.DistributedLockProvider = _ => new NoopDistributedSynchronizationProvider();
runtime.DistributedLockProvider = _ => new FileDistributedSynchronizationProvider(new DirectoryInfo(Path.Combine(Directory.GetCurrentDirectory(), "App_Data", "locks")));
runtime.WorkflowInboxCleanupOptions = options => configuration.GetSection("Runtime:WorkflowInboxCleanup").Bind(options);
runtime.WorkflowDispatcherOptions = options => configuration.GetSection("Runtime:WorkflowDispatcher").Bind(options);
})

View file

@ -41,7 +41,7 @@ public class WorkflowDefinitionEventsConsumer(IActivityRegistryPopulator activit
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionRetracted> context)
{
activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), context.Message.Id);
return Task.CompletedTask;
}
@ -79,7 +79,7 @@ public class WorkflowDefinitionEventsConsumer(IActivityRegistryPopulator activit
if (usableAsActivity)
return activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), id);
activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
return Task.CompletedTask;
}
}

View file

@ -30,7 +30,7 @@ public class RefreshActivityRegistry(IActivityRegistryPopulator activityRegistry
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken)
{
activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), notification.WorkflowDefinition.Id, cancellationToken);
activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), notification.WorkflowDefinition.Id, cancellationToken);
return Task.CompletedTask;
}
@ -61,7 +61,7 @@ public class RefreshActivityRegistry(IActivityRegistryPopulator activityRegistry
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionVersionDeleted notification, CancellationToken cancellationToken)
{
activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), notification.WorkflowDefinition.Id, cancellationToken);
activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), notification.WorkflowDefinition.Id, cancellationToken);
return Task.CompletedTask;
}
@ -81,7 +81,7 @@ public class RefreshActivityRegistry(IActivityRegistryPopulator activityRegistry
if (usableAsActivity.GetValueOrDefault())
return activityRegistryPopulator.AddToRegistry(typeof(WorkflowDefinitionActivityProvider), id);
activityRegistryPopulator.RemoveDefinitionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
activityRegistryPopulator.RemoveDefinitionVersionFromRegistry(typeof(WorkflowDefinitionActivityProvider), id);
return Task.CompletedTask;
}
}

View file

@ -37,7 +37,7 @@ public class ActivityRegistryPopulator(IEnumerable<IActivityProvider> providers,
var descriptorsToRemove = providerDescriptors
.Where(d =>
d.CustomProperties.TryGetValue("WorkflowDefinitionId", out var val) &&
val.ToString() == workflowDefinitionId);
val.ToString() == workflowDefinitionId).ToList();
foreach (ActivityDescriptor activityDescriptor in descriptorsToRemove)
{

View file

@ -0,0 +1,15 @@
using Elsa.MassTransit.Messages;
using Hangfire.Annotations;
using MassTransit;
namespace Elsa.Workflows.ComponentTests.Consumers;
[UsedImplicitly]
public class WorkflowDefinitionEventHandlers(IWorkflowDefinitionEvents workflowDefinitionEvents) : IConsumer<WorkflowDefinitionDeleted>
{
public Task Consume(ConsumeContext<WorkflowDefinitionDeleted> context)
{
workflowDefinitionEvents.OnWorkflowDefinitionDeleted(new WorkflowDefinitionDeletedEventArgs(context.Message.Id));
return Task.CompletedTask;
}
}

View file

@ -0,0 +1,8 @@
namespace Elsa.Workflows.ComponentTests;
public interface IWorkflowDefinitionEvents
{
event EventHandler<WorkflowDefinitionDeletedEventArgs> WorkflowDefinitionDeleted;
void OnWorkflowDefinitionDeleted(WorkflowDefinitionDeletedEventArgs args);
}

View file

@ -0,0 +1,6 @@
namespace Elsa.Workflows.ComponentTests;
public class WorkflowDefinitionDeletedEventArgs(string definitionId) : EventArgs
{
public string DefinitionId { get; } = definitionId;
}

View file

@ -1,9 +1,11 @@
using System.Net.Http.Headers;
using System.Reflection;
using Elsa.EntityFrameworkCore.Extensions;
using Elsa.EntityFrameworkCore.Modules.Management;
using Elsa.Extensions;
using Elsa.Identity.Providers;
using Elsa.MassTransit.Extensions;
using Elsa.Workflows.ComponentTests.Consumers;
using Elsa.Workflows.ComponentTests.Services;
using FluentStorage;
using Hangfire.Annotations;
@ -52,7 +54,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
elsa.UseDefaultAuthentication(defaultAuthentication => defaultAuthentication.UseAdminApiKey());
elsa.UseFluentStorageProvider(sp =>
{
var assemblyLocation = System.Reflection.Assembly.GetExecutingAssembly().Location;
var assemblyLocation = Assembly.GetExecutingAssembly().Location;
var assemblyDirectory = Path.GetDirectoryName(assemblyLocation)!;
var workflowsDirectorySegments = new[]
{
@ -64,6 +66,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
elsa.UseMassTransit(massTransit =>
{
massTransit.UseRabbitMq(rabbitMqConnectionString);
massTransit.AddConsumer<WorkflowDefinitionEventHandlers>("elsa-test-workflow-definition-updates", true);
});
elsa.UseWorkflowManagement(management =>
{
@ -82,6 +85,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
{
services.AddSingleton<ISignalManager, SignalManager>();
services.AddSingleton<IWorkflowEvents, WorkflowEvents>();
services.AddSingleton<IWorkflowDefinitionEvents, WorkflowDefinitionEvents>();
services.AddNotificationHandlersFrom<WorkflowServer>();
});
}

View file

@ -0,0 +1,7 @@
namespace Elsa.Workflows.ComponentTests.Services;
public class WorkflowDefinitionEvents : IWorkflowDefinitionEvents
{
public event EventHandler<WorkflowDefinitionDeletedEventArgs>? WorkflowDefinitionDeleted;
public void OnWorkflowDefinitionDeleted(WorkflowDefinitionDeletedEventArgs args) => WorkflowDefinitionDeleted?.Invoke(this, args);
}

View file

@ -0,0 +1,86 @@
using Elsa.Workflows.ComponentTests.Scenarios.WorkflowActivities.Workflows;
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
using Elsa.Workflows.Management.Contracts;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowActivities;
public class DeleteWorkflowTests : AppComponentTest
{
private readonly ISignalManager _signalManager;
private readonly IWorkflowDefinitionEvents _workflowDefinitionEvents;
private readonly IServiceScope _scope1;
private readonly IServiceScope _scope2;
private readonly IServiceScope _scope3;
private static readonly object WorkflowDeletedSignal = new();
public DeleteWorkflowTests(App app) : base(app)
{
_scope1 = app.Cluster.Pod1.Services.CreateScope();
_scope2 = app.Cluster.Pod2.Services.CreateScope();
_scope3 = app.Cluster.Pod3.Services.CreateScope();
_signalManager = Scope.ServiceProvider.GetRequiredService<ISignalManager>();
_workflowDefinitionEvents = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionEvents>();
_workflowDefinitionEvents.WorkflowDefinitionDeleted += OnWorkflowDefinionDeleted;
}
[Fact]
public async Task DeleteWorkflow()
{
EnsureWorkflowInRegistry(_scope1);
var workflowDefinitionManager = _scope1.ServiceProvider.GetRequiredService<IWorkflowDefinitionManager>();
await workflowDefinitionManager.DeleteByDefinitionIdAsync(Workflows.DeleteWorkflow.DefinitionId);
WorkflowTypeDeletedFromRegistry(_scope1);
}
[Fact]
public async Task DeleteWorkflow_Clustered()
{
EnsureWorkflowInRegistry(_scope1);
EnsureWorkflowInRegistry(_scope2);
EnsureWorkflowInRegistry(_scope3);
var workflowDefinitionManager = _scope3.ServiceProvider.GetRequiredService<IWorkflowDefinitionManager>();
await workflowDefinitionManager.DeleteByDefinitionIdAsync(Workflows.DeleteWorkflow.DefinitionId);
WorkflowTypeDeletedFromRegistry(_scope1);
await _signalManager.WaitAsync<WorkflowDefinitionDeletedEventArgs>(WorkflowDeletedSignal);
WorkflowTypeDeletedFromRegistry(_scope2);
WorkflowTypeDeletedFromRegistry(_scope3);
}
private void EnsureWorkflowInRegistry(IServiceScope scope)
{
var activityRegistry = scope.ServiceProvider.GetRequiredService<IActivityRegistry>();
var descriptor = activityRegistry.Find(Workflows.DeleteWorkflow.Type);
if (descriptor is null)
activityRegistry.Add(typeof(WorkflowDefinitionActivityProvider), descriptor);
}
private void WorkflowTypeDeletedFromRegistry(IServiceScope scope)
{
var activityRegistry = scope.ServiceProvider.GetRequiredService<IActivityRegistry>();
var descriptor = activityRegistry.Find(Workflows.DeleteWorkflow.Type);
Assert.Null(descriptor);
}
private void OnWorkflowDefinionDeleted(object? sender, WorkflowDefinitionDeletedEventArgs args)
{
if (args.DefinitionId == Workflows.DeleteWorkflow.DefinitionId)
{
_signalManager.Trigger(WorkflowDeletedSignal, args);
}
}
protected override void OnDispose()
{
_scope1.Dispose();
_scope2.Dispose();
_scope3.Dispose();
}
}

View file

@ -0,0 +1,61 @@
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Management.Models;
using Elsa.Workflows.Models;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowActivities;
public class SaveWorkflowTests(App app) : AppComponentTest(app)
{
private readonly IServiceScope _scope = app.Cluster.Pod1.Services.CreateScope();
[Theory]
[InlineData("Save1", true, true, true, true)]
[InlineData("Save2", true, false, true, false)]
[InlineData("Save3", false, true, false, false)]
[InlineData("Save4", false, false, false, false)]
private async Task SaveWorkflow(string name, bool usableAsActivity, bool publish, bool expectedInRegistry, bool isBrowsable)
{
var activityRegistry = _scope.ServiceProvider.GetRequiredService<IActivityRegistry>();
var descriptor = activityRegistry.Find(name);
if (descriptor is not null)
activityRegistry.Remove(typeof(WorkflowDefinitionActivityProvider), descriptor);
var importer = _scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionImporter>();
var request = new SaveWorkflowDefinitionRequest
{
Model = new WorkflowDefinitionModel
{
Name = name,
DefinitionId = name,
Options = new WorkflowOptions
{
UsableAsActivity = usableAsActivity,
AutoUpdateConsumingWorkflows = true
}
},
Publish = publish
};
await importer.ImportAsync(request);
descriptor = activityRegistry.Find(name);
if (expectedInRegistry)
{
Assert.NotNull(descriptor);
Assert.Equal(isBrowsable, descriptor.IsBrowsable);
}
else
{
Assert.Null(descriptor);
}
}
protected override void OnDispose()
{
_scope.Dispose();
}
}

View file

@ -0,0 +1,25 @@
using Elsa.Scheduling.Activities;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowActivities.Workflows;
public class DeleteWorkflow : WorkflowBase
{
public static readonly string DefinitionId = Guid.NewGuid().ToString();
public static readonly string Type = nameof(DeleteWorkflow);
protected override void Build(IWorkflowBuilder builder)
{
builder.Name = Type;
builder.WithDefinitionId(DefinitionId);
builder.WorkflowOptions.UsableAsActivity = true;
builder.Root = new Sequence
{
Activities =
{
new Delay(TimeSpan.FromMilliseconds(250)),
new WriteLine("This workflow will be deleted!")
}
};
}
}