From f448e9520a025613092ae1d73ea7e3af50fd28cd Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 17 Jul 2024 09:51:45 +0200 Subject: [PATCH] Fix reloading logic for workflow definitions (#5781) * Add reloading logic for workflow definitions Introduced a new mechanism to handle reloaded workflow definitions. This includes creating a "ReloadedWorkflowDefinition" model, updating the caching logic, and modifying notification handlers to work with the enhanced workflow reloading logic. This ensures workflow definitions are updated and managed correctly when published, retracted, or deleted. * Improve RefreshActivityRegistry documentation Updated the XML documentation to clarify that `RefreshActivityRegistry` refreshes the `IActivityRegistry` for `WorkflowDefinitionActivityProvider` whenever workflow definitions are reloaded, instead of when they are published, retracted, or deleted. * Add 'materializer' to user dictionary The term 'materializer' has been added to the user dictionary to improve code spelling and naming consistency. This change ensures that 'materializer' is recognized as a correct term in the codebase. * Organize test files by adding a 'Workflows' directory Renamed 'http-workflow.json' to indicate it belongs under 'Workflows'. This improves file organization and clarity within the 'WorkflowDefinitionReload' scenario. * Add TestWorkflowMaterializer and TestWorkflowProvider Introduce `TestWorkflowMaterializer` for deserializing workflows from `TestWorkflowProvider`. Added integration of these new components in the `ReloadWorkflowTests` and `WorkflowServer`. Also renamed `RemoveReloadWorkflowTests` to `ReloadWorkflowTests`. * Rename and expand workflow reload tests Renamed `RemoveReloadWorkflowTests` to `ReloadWorkflowTests` to better reflect its purpose and expanded with additional test cases. Added tests to verify workflow and activity registry updates after source provider changes and workflow reloads. * Update copy settings for workflow test scenarios Reorganized and added 'CopyToOutputDirectory' settings for JSON files in workflow test scenarios. Ensured all necessary files are correctly included and copied during output directory builds to maintain test consistency. --- Elsa.sln.DotSettings | 2 + .../Handlers/InvalidateHttpWorkflowsCache.cs | 10 +- .../WorkflowDefinitionEventsConsumer.cs | 2 +- ...dWorkflowDefinitionNotificationsHandler.cs | 4 +- .../Messages/WorkflowDefinitionsReloaded.cs | 10 +- ...rkflowDefinitionActivityRegistryUpdater.cs | 5 +- .../Handlers/DeleteWorkflowInstances.cs | 2 + .../IWorkflowDefinitionStorePopulator.cs | 13 +- .../Contracts/IWorkflowDefinitionsReloader.cs | 4 +- .../Contracts/IWorkflowRuntime.cs | 2 +- .../Features/CachingWorkflowRuntimeFeature.cs | 3 +- .../Features/WorkflowRuntimeFeature.cs | 3 +- .../Handlers/InvalidateWorkflowsCache.cs | 8 +- .../Handlers/RefreshActivityRegistry.cs | 31 +++++ .../Models/ReloadedWorkflowDefinition.cs | 30 +++++ .../WorkflowDefinitionsReloaded.cs | 3 +- ...DefaultWorkflowDefinitionStorePopulator.cs | 35 +++--- .../Services/WorkflowDefinitionsReloader.cs | 6 +- .../Elsa.Workflows.ComponentTests.csproj | 6 +- .../Helpers/Fixtures/WorkflowServer.cs | 17 ++- .../Materializers/TestWorkflowMaterializer.cs | 26 ++++ .../WorkflowProviders/TestWorkflowProvider.cs | 14 +++ .../ReloadWorkflowTests.cs | 113 ++++++++++++++++++ .../RemoveReloadWorkflowTests.cs | 35 ------ .../{ => Workflows}/http-workflow.json | 0 25 files changed, 290 insertions(+), 94 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs delete mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs rename test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/{ => Workflows}/http-workflow.json (100%) diff --git a/Elsa.sln.DotSettings b/Elsa.sln.DotSettings index 19414abf6..b291186ee 100644 --- a/Elsa.sln.DotSettings +++ b/Elsa.sln.DotSettings @@ -12,10 +12,12 @@ EF True True + True True True True True + True True True True diff --git a/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs b/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs index bb87afe99..f07a32072 100644 --- a/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs +++ b/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs @@ -95,10 +95,8 @@ public class InvalidateHttpWorkflowsCache( /// public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) { - foreach (var workflowDefinitionId in notification.WorkflowDefinitionIds) - { - await InvalidateCacheAsync(workflowDefinitionId); - } + foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions) + await InvalidateCacheAsync(reloadedWorkflowDefinition.DefinitionId); } private async Task InvalidateCacheAsync(string workflowDefinitionId) @@ -119,9 +117,9 @@ public class InvalidateHttpWorkflowsCache( private async Task InvalidateTriggerCacheAsync(IEnumerable triggers, CancellationToken cancellationToken) { - foreach (StoredTrigger trigger in triggers) + foreach (var trigger in triggers) { - if (trigger?.Payload is HttpEndpointBookmarkPayload httpPayload) + if (trigger.Payload is HttpEndpointBookmarkPayload httpPayload) { var hash = httpWorkflowsCacheManager.ComputeBookmarkHash(httpPayload.Path, httpPayload.Method); await httpWorkflowsCacheManager.EvictTriggerAsync(hash, cancellationToken); diff --git a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs index d43d11172..4f7821a63 100644 --- a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs +++ b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs @@ -92,7 +92,7 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr public async Task Consume(ConsumeContext context) { var message = context.Message; - var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsReloaded(message.WorkflowDefinitionIds); + var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsReloaded(message.ReloadedWorkflowDefinitions); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = true; await notificationSender.SendAsync(notification, context.CancellationToken); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = false; diff --git a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs index 991eb0914..10e654479 100644 --- a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs +++ b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs @@ -88,8 +88,8 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IBus bus) : if (AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer) return Task.CompletedTask; - var definitionIds = notification.WorkflowDefinitionIds; - var message = new Distributed.WorkflowDefinitionsReloaded(definitionIds); + var reloadedWorkflowDefinitions = notification.ReloadedWorkflowDefinitions; + var message = new Distributed.WorkflowDefinitionsReloaded(reloadedWorkflowDefinitions); return bus.Publish(message, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs index bba459390..138492683 100644 --- a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs +++ b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs @@ -1,8 +1,10 @@ -namespace Elsa.MassTransit.Messages; +using Elsa.Workflows.Runtime.Models; + +namespace Elsa.MassTransit.Messages; /// Represents a message that indicates that the specified workflow definitions have been reloaded. -public class WorkflowDefinitionsReloaded(ICollection workflowDefinitionIds) +public class WorkflowDefinitionsReloaded(ICollection reloadedWorkflowDefinitions) { - /// The workflow definition IDs that have been reloaded. - public ICollection WorkflowDefinitionIds { get; set; } = workflowDefinitionIds; + /// The reloaded workflow definitions. + public ICollection ReloadedWorkflowDefinitions { get; set; } = reloadedWorkflowDefinitions; } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs index 454ac3dbe..826b37c8c 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs @@ -8,9 +8,9 @@ public interface IWorkflowDefinitionActivityRegistryUpdater /// /// Tries to add a workflow as an activity to the registry. /// - /// The ID of the workflow definition. + /// The version ID of the workflow definition. /// The cancellation token. - Task AddToRegistry(string workflowDefinitionId, CancellationToken cancellationToken = default); + Task AddToRegistry(string workflowDefinitionVersionId, CancellationToken cancellationToken = default); /// /// Removes workflow definition activities from the . @@ -18,7 +18,6 @@ public interface IWorkflowDefinitionActivityRegistryUpdater /// The ID of the workflow definition to remove. void RemoveDefinitionFromRegistry(string workflowDefinitionId); - /// /// Removes a workflow definition version activity from the . /// diff --git a/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs b/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs index 1a0178476..66141f945 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs @@ -2,12 +2,14 @@ using Elsa.Mediator.Contracts; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Notifications; +using JetBrains.Annotations; namespace Elsa.Workflows.Management.Handlers; /// /// Deletes workflow instances when a workflow definition or version is deleted. /// +[UsedImplicitly] public class DeleteWorkflowInstances : INotificationHandler, INotificationHandler, diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs index aeb5c198b..5cdf956dc 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs @@ -1,4 +1,5 @@ using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Runtime.Models; namespace Elsa.Workflows.Runtime.Contracts; @@ -12,27 +13,27 @@ public interface IWorkflowDefinitionStorePopulator /// Populates the with workflow definitions provided from implementations. /// /// The cancellation token. - Task> PopulateStoreAsync(CancellationToken cancellationToken = default); + Task> PopulateStoreAsync(CancellationToken cancellationToken = default); /// /// Populates the with workflow definitions provided from implementations. /// /// Whether to index triggers. /// The cancellation token. - Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default); + Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default); /// /// Adds a workflow definition to the store. /// /// A materialized workflow. /// An optional cancellation token. - Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default); - + Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default); + /// /// Adds a workflow definition to the store. /// /// A materialized workflow. - /// /// Whether to index triggers. + /// Whether to index triggers. /// An optional cancellation token. - Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default); + Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs index e2cc2987e..f17d06edf 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs @@ -1,8 +1,8 @@ namespace Elsa.Workflows.Runtime.Contracts; -/// Reloads all workflows by re-invoking the populator. +/// Reloads all workflows by invoking the populator. public interface IWorkflowDefinitionsReloader { - /// Reloads all workflows by re-invoking the populator. + /// Reloads all workflows by invoking the populator. Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs index a60d31f88..be218983b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs @@ -15,7 +15,7 @@ namespace Elsa.Workflows.Runtime.Contracts; public interface IWorkflowRuntime { /// - /// Returns a value whether or not the specified workflow definition can create a new instance. + /// Returns a value whether the specified workflow definition can create a new instance. /// Task CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = default); diff --git a/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs index 9c0e26e73..bb31de5d4 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs @@ -26,6 +26,7 @@ public class CachingWorkflowRuntimeFeature : FeatureBase .Decorate() // Handlers. - .AddNotificationHandler(); + .AddNotificationHandler() + .AddNotificationHandler(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index 75a41922e..6bdbb7d99 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -7,7 +7,6 @@ using Elsa.Features.Attributes; using Elsa.Features.Services; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; -using Elsa.Workflows.Management.Handlers; using Elsa.Workflows.Management.Services; using Elsa.Workflows.Runtime.ActivationValidators; using Elsa.Workflows.Runtime.Contracts; @@ -272,12 +271,12 @@ public class WorkflowRuntimeFeature : FeatureBase .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() - .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() // Workflow activation strategies. .AddScoped() diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs index fb6dc6d45..d2d52613c 100644 --- a/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs @@ -6,7 +6,7 @@ using JetBrains.Annotations; namespace Elsa.Workflows.Runtime.Handlers; /// -/// A notification handler that invalidates workflows cache when workflow definitions are reloaded. +/// A notification handler that invalidates the workflow cache when workflow definitions are reloaded. /// /// /// The class implements the INotificationHandler interface and is responsible for handling WorkflowDefinitionsReloaded notifications. @@ -18,9 +18,7 @@ public class InvalidateWorkflowsCache(IWorkflowDefinitionCacheManager workflowDe /// public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) { - foreach (var workflowDefinitionId in notification.WorkflowDefinitionIds) - { - await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(workflowDefinitionId, cancellationToken); - } + foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions) + await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(reloadedWorkflowDefinition.DefinitionId, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs new file mode 100644 index 000000000..607751270 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs @@ -0,0 +1,31 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Runtime.Notifications; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Handlers; + +/// Refreshes the for the provider whenever workflow definitions are reloaded. +[PublicAPI] +public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater workflowDefinitionActivityRegistryUpdater) : INotificationHandler +{ + /// + public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) + { + foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions) + await UpdateDefinition(reloadedWorkflowDefinition.DefinitionVersionId, reloadedWorkflowDefinition.UsableAsActivity); + } + + private Task UpdateDefinition(string definitionVersionId, bool? usableAsActivity) + { + // A workflow should remain in the activity registry unless no longer being marked as an activity. + if (usableAsActivity.GetValueOrDefault()) + return workflowDefinitionActivityRegistryUpdater.AddToRegistry(definitionVersionId); + + workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(definitionVersionId); + return Task.CompletedTask; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs b/src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs new file mode 100644 index 000000000..7a30f1f3b --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs @@ -0,0 +1,30 @@ +using Elsa.Workflows.Management.Entities; + +namespace Elsa.Workflows.Runtime.Models; + +/// +/// Represents a reloaded workflow definition with necessary properties for +/// identification, versioning, and usability status as an activity. +/// +/// The unique identifier for the workflow definition. +/// The unique identifier for the specific version of the workflow definition. +/// The version number of the workflow definition. +/// Indicates whether the workflow definition can be used as an activity. +public record ReloadedWorkflowDefinition(string DefinitionId, string DefinitionVersionId, int Version, bool UsableAsActivity) +{ + /// + /// Creates an instance of from a given . + /// + /// The workflow definition used to create the reloaded workflow definition. + /// A new instance of . + public static ReloadedWorkflowDefinition FromDefinition(WorkflowDefinition workflowDefinition) + { + return new ReloadedWorkflowDefinition + ( + workflowDefinition.DefinitionId, + workflowDefinition.Id, + workflowDefinition.Version, + workflowDefinition.Options.UsableAsActivity ?? false + ); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs index 882fa7a51..336347ff4 100644 --- a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs @@ -1,6 +1,7 @@ using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Models; namespace Elsa.Workflows.Runtime.Notifications; /// Published when workflow definitions have been reloaded. -public record WorkflowDefinitionsReloaded(ICollection WorkflowDefinitionIds) : INotification; +public record WorkflowDefinitionsReloaded(ICollection ReloadedWorkflowDefinitions) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs index cba487b5a..fad0f8759 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -49,42 +49,47 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } /// - public Task> PopulateStoreAsync(CancellationToken cancellationToken = default) + public Task> PopulateStoreAsync(CancellationToken cancellationToken = default) { return PopulateStoreAsync(true, cancellationToken); } /// - public async Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default) + public async Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default) { var providers = _workflowDefinitionProviders(); - var workflowDefinitionIds = new List(); + var workflowDefinitions = new List(); + foreach (var provider in providers) { var results = await provider.GetWorkflowsAsync(cancellationToken).AsTask().ToList(); - workflowDefinitionIds.AddRange(results.Select(w => w.Workflow.Id)); - - foreach (var result in results) await AddAsync(result, indexTriggers, cancellationToken); + foreach (var result in results) + { + var workflowDefinition = await AddAsync(result, indexTriggers, cancellationToken); + workflowDefinitions.Add(workflowDefinition); + } } - return workflowDefinitionIds; + return workflowDefinitions; } /// - public Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) + public Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) { return AddAsync(materializedWorkflow, true, cancellationToken); } /// - public async Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default) + public async Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default) { await AssignIdentities(materializedWorkflow.Workflow, cancellationToken); - await AddOrUpdateAsync(materializedWorkflow, cancellationToken); + var workflowDefinition = await AddOrUpdateAsync(materializedWorkflow, cancellationToken); if (indexTriggers) await IndexTriggersAsync(materializedWorkflow, cancellationToken); + + return workflowDefinition; } private async Task AssignIdentities(Workflow workflow, CancellationToken cancellationToken) @@ -92,13 +97,13 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP await _identityGraphService.AssignIdentitiesAsync(workflow, cancellationToken); } - private async Task AddOrUpdateAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) + private async Task AddOrUpdateAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) { await _semaphore.WaitAsync(cancellationToken); try { - await AddOrUpdateCoreAsync(materializedWorkflow, cancellationToken); + return await AddOrUpdateCoreAsync(materializedWorkflow, cancellationToken); } finally { @@ -106,7 +111,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } } - private async Task AddOrUpdateCoreAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) + private async Task AddOrUpdateCoreAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) { var workflow = materializedWorkflow.Workflow; var definitionId = workflow.Identity.DefinitionId; @@ -178,7 +183,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP if (existingDefinitionVersion is null && workflowDefinitionsToSave.Any(w => w.Id == workflowDefinition.Id)) { _logger.LogInformation("Workflow with ID {WorkflowId} already exists", workflowDefinition.Id); - return; + return workflowDefinition; } workflowDefinitionsToSave.Add(workflowDefinition); @@ -194,7 +199,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } await _workflowDefinitionStore.SaveManyAsync(workflowDefinitionsToSave, cancellationToken); - return; + return workflowDefinition; async Task UpdateIsLatest() { diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs index 3afadafd2..f1a179e51 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs @@ -1,5 +1,6 @@ using Elsa.Mediator.Contracts; using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Notifications; namespace Elsa.Workflows.Runtime.Services; @@ -10,8 +11,9 @@ public class WorkflowDefinitionsReloader(IWorkflowDefinitionStorePopulator workf /// public async Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken) { - var definitionIds = await workflowDefinitionStorePopulator.PopulateStoreAsync(true, cancellationToken); - var notification = new WorkflowDefinitionsReloaded(definitionIds); + var workflowDefinitions = await workflowDefinitionStorePopulator.PopulateStoreAsync(true, cancellationToken); + var reloadedWorkflowDefinitions = workflowDefinitions.Select(ReloadedWorkflowDefinition.FromDefinition).ToList(); + var notification = new WorkflowDefinitionsReloaded(reloadedWorkflowDefinitions); await notificationSender.SendAsync(notification, cancellationToken); } } \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj index c6588808f..deb69e04d 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj +++ b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj @@ -21,9 +21,6 @@ - - Always - Always @@ -81,6 +78,9 @@ Always + + Always + diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs index aaaa665b4..a026a35a4 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs @@ -6,8 +6,11 @@ using Elsa.Extensions; using Elsa.Identity.Providers; using Elsa.MassTransit.Extensions; using Elsa.Workflows.ComponentTests.Consumers; +using Elsa.Workflows.ComponentTests.Helpers.Materializers; using Elsa.Workflows.ComponentTests.Helpers.Services; +using Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; using Elsa.Workflows.ComponentTests.Services; +using Elsa.Workflows.Management.Contracts; using FluentStorage; using Hangfire.Annotations; using Microsoft.AspNetCore.Hosting; @@ -94,11 +97,15 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl builder.ConfigureTestServices(services => { - services.AddSingleton(); - services.AddSingleton(); - services.AddSingleton(); - services.AddSingleton(); - services.AddNotificationHandlersFrom(); + services + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddScoped() + .AddNotificationHandlersFrom() + .AddWorkflowDefinitionProvider() + ; }); } diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs new file mode 100644 index 000000000..e011422c4 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs @@ -0,0 +1,26 @@ +using Elsa.Workflows.Activities; +using Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Runtime.Contracts; + +namespace Elsa.Workflows.ComponentTests.Helpers.Materializers; + +/// A workflow materializer that deserializes workflows created from . +public class TestWorkflowMaterializer(IEnumerable workflowProviders) : IWorkflowMaterializer +{ + /// The name of the materializer. + public const string MaterializerName = "Test"; + + /// + public string Name => MaterializerName; + + /// + public ValueTask MaterializeAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + { + var testProvider = (TestWorkflowProvider)workflowProviders.Single(x => x is TestWorkflowProvider); + var materializedWorkflow = testProvider.MaterializedWorkflows.First(x => x.Workflow.Identity.Id == definition.Id); + + return ValueTask.FromResult(materializedWorkflow.Workflow); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs new file mode 100644 index 000000000..5473cf8eb --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs @@ -0,0 +1,14 @@ +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; + +namespace Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; + +public class TestWorkflowProvider : IWorkflowProvider +{ + public string Name => "Test"; + public ICollection MaterializedWorkflows { get; set; } = new List(); + public ValueTask> GetWorkflowsAsync(CancellationToken cancellationToken = default) + { + return new(MaterializedWorkflows); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs new file mode 100644 index 000000000..120f0fe42 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -0,0 +1,113 @@ +using System.Net; +using Elsa.Common.Models; +using Elsa.Workflows.Activities; +using Elsa.Workflows.ComponentTests.Helpers.Materializers; +using Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Materializers; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowDefinitionReload; + +public class ReloadWorkflowTests : AppComponentTest +{ + private readonly IWorkflowDefinitionManager _workflowDefinitionManager; + private readonly IWorkflowDefinitionsReloader _workflowDefinitionsReloader; + private readonly IWorkflowBuilderFactory _workflowBuilderFactory; + private readonly TestWorkflowProvider _testWorkflowProvider; + private readonly IWorkflowDefinitionService _workflowDefinitionService; + private readonly IActivityRegistry _activityRegistry; + + public ReloadWorkflowTests(App app) : base(app) + { + _workflowDefinitionManager = Scope.ServiceProvider.GetRequiredService(); + _workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService(); + _workflowBuilderFactory = Scope.ServiceProvider.GetRequiredService(); + _workflowDefinitionService = Scope.ServiceProvider.GetRequiredService(); + _activityRegistry = Scope.ServiceProvider.GetRequiredService(); + var workflowProviders = Scope.ServiceProvider.GetRequiredService>(); + _testWorkflowProvider = (TestWorkflowProvider)workflowProviders.First(x => x is TestWorkflowProvider); + } + + [Fact] + public async Task Reloading_AfterRemovingTheWorkflow_ShouldMakeWorkflowReachableAgain() + { + var client = WorkflowServer.CreateHttpWorkflowClient(); + await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None); + var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None); + var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); + Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); + } + + [Fact] + public async Task Reloading_AfterUpdatingSourceProvider_ShouldRefreshCaches() + { + var definitionId = Guid.NewGuid().ToString(); + var definitionVersionId1 = Guid.NewGuid().ToString(); + var workflowV1 = await BuildWorkflowAsync(definitionId, definitionVersionId1, 1); + + // Set up the initial workflow version. + _testWorkflowProvider.MaterializedWorkflows = [workflowV1]; + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + var definitionV1 = await _workflowDefinitionService.FindWorkflowGraphAsync(definitionId, VersionOptions.Latest); + Assert.Equal(definitionVersionId1, definitionV1!.Workflow.Identity.Id); + + // Simulate the workflow provider to have a new version available. + var definitionVersionId2 = Guid.NewGuid().ToString(); + var workflowV2 = await BuildWorkflowAsync(definitionId, definitionVersionId2, 2); + _testWorkflowProvider.MaterializedWorkflows = [workflowV1, workflowV2]; + + // Reload the workflow definitions. + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + + // Assert that the workflow definition service finds the updated workflow version. + var definitionV2 = await _workflowDefinitionService.FindWorkflowGraphAsync(definitionId, VersionOptions.Latest); + Assert.Equal(definitionVersionId2, definitionV2!.Workflow.Identity.Id); + } + + [Fact] + public async Task Reloading_AfterUpdatingSourceProvider_ShouldRefreshActivityRegistry() + { + var definitionId = Guid.NewGuid().ToString(); + var definitionVersionId1 = Guid.NewGuid().ToString(); + var workflowV1 = await BuildWorkflowAsync(definitionId, definitionVersionId1, 1); + + // Set up the initial workflow version. + _testWorkflowProvider.MaterializedWorkflows = [workflowV1]; + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + var activityV1 = _activityRegistry.Find(workflowV1.Workflow.Name!); + Assert.Equal(1, activityV1!.Version); + + // Simulate the workflow provider to have a new version available. + var definitionVersionId2 = Guid.NewGuid().ToString(); + var workflowV2 = await BuildWorkflowAsync(definitionId, definitionVersionId2, 2); + _testWorkflowProvider.MaterializedWorkflows = [workflowV1, workflowV2]; + + // Reload the workflow definitions. + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + + // Assert that the activity registry contains a new activity descriptor representing the new workflow version. + var activityV2 = _activityRegistry.Find(workflowV2.Workflow.Name!)!; + Assert.Equal(2, activityV2.Version); + } + + private async Task BuildWorkflowAsync(string definitionId, string definitionVersionId, int version) + { + var builder = _workflowBuilderFactory.CreateBuilder(); + builder.DefinitionId = definitionId; + builder.Id = definitionVersionId; + builder.Version = version; + builder.Name = Guid.NewGuid().ToString(); + builder.Root = new WriteLine($"Version {version}"); + builder.WorkflowOptions.UsableAsActivity = true; + var workflow = await builder.BuildWorkflowAsync(); + workflow.Name = builder.Name; + return new MaterializedWorkflow(workflow, _testWorkflowProvider.Name, TestWorkflowMaterializer.MaterializerName); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs deleted file mode 100644 index 390c58d14..000000000 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs +++ /dev/null @@ -1,35 +0,0 @@ -using System.Net; -using Elsa.Workflows.Management.Contracts; -using Elsa.Workflows.Runtime.Contracts; -using Microsoft.Extensions.DependencyInjection; - -namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowDefinitionReload; - -public class RemoveReloadWorkflowTests : AppComponentTest -{ - private readonly IWorkflowDefinitionManager _workflowDefinitionManager; - private readonly IWorkflowDefinitionsReloader _workflowDefinitionsReloader; - - public RemoveReloadWorkflowTests(App app) : base(app) - { - _workflowDefinitionManager = Scope.ServiceProvider.GetRequiredService(); - _workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService(); - } - - [Fact] - public async Task RemovingTheWorkflowThenReload_WorkflowShouldBeReachableAgain() - { - var client = WorkflowServer.CreateHttpWorkflowClient(); - - var result = await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None); - - var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); - - await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None); - - var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); - - Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); - Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); - } -} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/Workflows/http-workflow.json similarity index 100% rename from test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json rename to test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/Workflows/http-workflow.json