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.
This commit is contained in:
parent
ab8deafcc1
commit
f448e9520a
|
|
@ -12,10 +12,12 @@
|
|||
<s:String x:Key="/Default/CodeStyle/Naming/CSharpNaming/Abbreviations/=EF/@EntryIndexedValue">EF</s:String>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=downloadables/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=initializable/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=materializer/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=materializers/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Persister/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Populator/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Postgre/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=reloader/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=startable/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Telnyx/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Unschedule/@EntryIndexedValue">True</s:Boolean>
|
||||
|
|
|
|||
|
|
@ -95,10 +95,8 @@ public class InvalidateHttpWorkflowsCache(
|
|||
/// <inheritdoc />
|
||||
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<StoredTrigger> 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);
|
||||
|
|
|
|||
|
|
@ -92,7 +92,7 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr
|
|||
public async Task Consume(ConsumeContext<WorkflowDefinitionsReloaded> 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;
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<string> workflowDefinitionIds)
|
||||
public class WorkflowDefinitionsReloaded(ICollection<ReloadedWorkflowDefinition> reloadedWorkflowDefinitions)
|
||||
{
|
||||
/// The workflow definition IDs that have been reloaded.
|
||||
public ICollection<string> WorkflowDefinitionIds { get; set; } = workflowDefinitionIds;
|
||||
/// The reloaded workflow definitions.
|
||||
public ICollection<ReloadedWorkflowDefinition> ReloadedWorkflowDefinitions { get; set; } = reloadedWorkflowDefinitions;
|
||||
}
|
||||
|
|
@ -8,9 +8,9 @@ public interface IWorkflowDefinitionActivityRegistryUpdater
|
|||
/// <summary>
|
||||
/// Tries to add a workflow as an activity to the registry.
|
||||
/// </summary>
|
||||
/// <param name="workflowDefinitionId">The ID of the workflow definition.</param>
|
||||
/// <param name="workflowDefinitionVersionId">The version ID of the workflow definition.</param>
|
||||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
Task AddToRegistry(string workflowDefinitionId, CancellationToken cancellationToken = default);
|
||||
Task AddToRegistry(string workflowDefinitionVersionId, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Removes workflow definition activities from the <see cref="Elsa.Workflows.Contracts.IActivityRegistry"/>.
|
||||
|
|
@ -18,7 +18,6 @@ public interface IWorkflowDefinitionActivityRegistryUpdater
|
|||
/// <param name="workflowDefinitionId">The ID of the workflow definition to remove.</param>
|
||||
void RemoveDefinitionFromRegistry(string workflowDefinitionId);
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Removes a workflow definition version activity from the <see cref="Elsa.Workflows.Contracts.IActivityRegistry"/>.
|
||||
/// </summary>
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Deletes workflow instances when a workflow definition or version is deleted.
|
||||
/// </summary>
|
||||
[UsedImplicitly]
|
||||
public class DeleteWorkflowInstances :
|
||||
INotificationHandler<WorkflowDefinitionDeleting>,
|
||||
INotificationHandler<WorkflowDefinitionVersionDeleting>,
|
||||
|
|
|
|||
|
|
@ -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 <see cref="IWorkflowDefinitionStore"/> with workflow definitions provided from <see cref="IWorkflowProvider"/> implementations.
|
||||
/// </summary>
|
||||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
Task<ICollection<string>> PopulateStoreAsync(CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Populates the <see cref="IWorkflowDefinitionStore"/> with workflow definitions provided from <see cref="IWorkflowProvider"/> implementations.
|
||||
/// </summary>
|
||||
/// <param name="indexTriggers">Whether to index triggers.</param>
|
||||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
Task<ICollection<string>> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Adds a workflow definition to the store.
|
||||
/// </summary>
|
||||
/// <param name="materializedWorkflow">A materialized workflow.</param>
|
||||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default);
|
||||
|
||||
Task<WorkflowDefinition> AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Adds a workflow definition to the store.
|
||||
/// </summary>
|
||||
/// <param name="materializedWorkflow">A materialized workflow.</param>
|
||||
/// /// <param name="indexTriggers">Whether to index triggers.</param>
|
||||
/// <param name="indexTriggers">Whether to index triggers.</param>
|
||||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default);
|
||||
Task<WorkflowDefinition> AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -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);
|
||||
}
|
||||
|
|
@ -15,7 +15,7 @@ namespace Elsa.Workflows.Runtime.Contracts;
|
|||
public interface IWorkflowRuntime
|
||||
{
|
||||
/// <summary>
|
||||
/// 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.
|
||||
/// </summary>
|
||||
Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = default);
|
||||
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ public class CachingWorkflowRuntimeFeature : FeatureBase
|
|||
.Decorate<ITriggerStore, CachingTriggerStore>()
|
||||
|
||||
// Handlers.
|
||||
.AddNotificationHandler<InvalidateTriggersCache>();
|
||||
.AddNotificationHandler<InvalidateTriggersCache>()
|
||||
.AddNotificationHandler<InvalidateWorkflowsCache>();
|
||||
}
|
||||
}
|
||||
|
|
@ -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<CancelBackgroundActivities>()
|
||||
.AddNotificationHandler<DeleteBookmarks>()
|
||||
.AddNotificationHandler<DeleteTriggers>()
|
||||
.AddNotificationHandler<DeleteWorkflowInstances>()
|
||||
.AddNotificationHandler<DeleteActivityExecutionLogRecords>()
|
||||
.AddNotificationHandler<ReadWorkflowInboxMessage>()
|
||||
.AddNotificationHandler<DeliverWorkflowMessagesFromInbox>()
|
||||
.AddNotificationHandler<DeleteWorkflowExecutionLogRecords>()
|
||||
.AddNotificationHandler<WorkflowExecutionContextNotificationsHandler>()
|
||||
.AddNotificationHandler<RefreshActivityRegistry>()
|
||||
|
||||
// Workflow activation strategies.
|
||||
.AddScoped<IWorkflowActivationStrategy, SingletonStrategy>()
|
||||
|
|
|
|||
|
|
@ -6,7 +6,7 @@ using JetBrains.Annotations;
|
|||
namespace Elsa.Workflows.Runtime.Handlers;
|
||||
|
||||
/// <summary>
|
||||
/// 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.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// The class implements the <c>INotificationHandler</c> interface and is responsible for handling <c>WorkflowDefinitionsReloaded</c> notifications.
|
||||
|
|
@ -18,9 +18,7 @@ public class InvalidateWorkflowsCache(IWorkflowDefinitionCacheManager workflowDe
|
|||
/// <inheritdoc />
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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 <see cref="IActivityRegistry"/> for the <see cref="WorkflowDefinitionActivityProvider"/> provider whenever workflow definitions are reloaded.
|
||||
[PublicAPI]
|
||||
public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater workflowDefinitionActivityRegistryUpdater) : INotificationHandler<WorkflowDefinitionsReloaded>
|
||||
{
|
||||
/// <inheritdoc />
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,30 @@
|
|||
using Elsa.Workflows.Management.Entities;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Models;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a reloaded workflow definition with necessary properties for
|
||||
/// identification, versioning, and usability status as an activity.
|
||||
/// </summary>
|
||||
/// <param name="DefinitionId">The unique identifier for the workflow definition.</param>
|
||||
/// <param name="DefinitionVersionId">The unique identifier for the specific version of the workflow definition.</param>
|
||||
/// <param name="Version">The version number of the workflow definition.</param>
|
||||
/// <param name="UsableAsActivity">Indicates whether the workflow definition can be used as an activity.</param>
|
||||
public record ReloadedWorkflowDefinition(string DefinitionId, string DefinitionVersionId, int Version, bool UsableAsActivity)
|
||||
{
|
||||
/// <summary>
|
||||
/// Creates an instance of <see cref="ReloadedWorkflowDefinition"/> from a given <see cref="WorkflowDefinition"/>.
|
||||
/// </summary>
|
||||
/// <param name="workflowDefinition">The workflow definition used to create the reloaded workflow definition.</param>
|
||||
/// <returns>A new instance of <see cref="ReloadedWorkflowDefinition"/>.</returns>
|
||||
public static ReloadedWorkflowDefinition FromDefinition(WorkflowDefinition workflowDefinition)
|
||||
{
|
||||
return new ReloadedWorkflowDefinition
|
||||
(
|
||||
workflowDefinition.DefinitionId,
|
||||
workflowDefinition.Id,
|
||||
workflowDefinition.Version,
|
||||
workflowDefinition.Options.UsableAsActivity ?? false
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<string> WorkflowDefinitionIds) : INotification;
|
||||
public record WorkflowDefinitionsReloaded(ICollection<ReloadedWorkflowDefinition> ReloadedWorkflowDefinitions) : INotification;
|
||||
|
|
@ -49,42 +49,47 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<ICollection<string>> PopulateStoreAsync(CancellationToken cancellationToken = default)
|
||||
public Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
return PopulateStoreAsync(true, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ICollection<string>> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default)
|
||||
public async Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var providers = _workflowDefinitionProviders();
|
||||
var workflowDefinitionIds = new List<string>();
|
||||
var workflowDefinitions = new List<WorkflowDefinition>();
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default)
|
||||
public Task<WorkflowDefinition> AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default)
|
||||
{
|
||||
return AddAsync(materializedWorkflow, true, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default)
|
||||
public async Task<WorkflowDefinition> 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<WorkflowDefinition> 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<WorkflowDefinition> 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()
|
||||
{
|
||||
|
|
|
|||
|
|
@ -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
|
|||
/// <inheritdoc />
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
@ -21,9 +21,6 @@
|
|||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<None Update="Scenarios\WorkflowDefinitionReload\http-workflow.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
<None Update="xunit.runner.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
|
|
@ -81,6 +78,9 @@
|
|||
<None Update="Scenarios\WorkflowDefinitionReload\http-workflow.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
<None Update="Scenarios\WorkflowDefinitionReload\Workflows\http-workflow.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
|
|||
|
|
@ -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<ISignalManager, SignalManager>();
|
||||
services.AddSingleton<IWorkflowEvents, WorkflowEvents>();
|
||||
services.AddSingleton<IWorkflowDefinitionEvents, WorkflowDefinitionEvents>();
|
||||
services.AddSingleton<ITriggerChangeTokenSignalEvents, TriggerChangeTokenSignalEvents>();
|
||||
services.AddNotificationHandlersFrom<WorkflowServer>();
|
||||
services
|
||||
.AddSingleton<ISignalManager, SignalManager>()
|
||||
.AddSingleton<IWorkflowEvents, WorkflowEvents>()
|
||||
.AddSingleton<IWorkflowDefinitionEvents, WorkflowDefinitionEvents>()
|
||||
.AddSingleton<ITriggerChangeTokenSignalEvents, TriggerChangeTokenSignalEvents>()
|
||||
.AddScoped<IWorkflowMaterializer, TestWorkflowMaterializer>()
|
||||
.AddNotificationHandlersFrom<WorkflowServer>()
|
||||
.AddWorkflowDefinitionProvider<TestWorkflowProvider>()
|
||||
;
|
||||
});
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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 <see cref="TestWorkflowProvider"/>.
|
||||
public class TestWorkflowMaterializer(IEnumerable<IWorkflowProvider> workflowProviders) : IWorkflowMaterializer
|
||||
{
|
||||
/// The name of the materializer.
|
||||
public const string MaterializerName = "Test";
|
||||
|
||||
/// <inheritdoc />
|
||||
public string Name => MaterializerName;
|
||||
|
||||
/// <inheritdoc />
|
||||
public ValueTask<Workflow> 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<MaterializedWorkflow> MaterializedWorkflows { get; set; } = new List<MaterializedWorkflow>();
|
||||
public ValueTask<IEnumerable<MaterializedWorkflow>> GetWorkflowsAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
return new(MaterializedWorkflows);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<IWorkflowDefinitionManager>();
|
||||
_workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionsReloader>();
|
||||
_workflowBuilderFactory = Scope.ServiceProvider.GetRequiredService<IWorkflowBuilderFactory>();
|
||||
_workflowDefinitionService = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionService>();
|
||||
_activityRegistry = Scope.ServiceProvider.GetRequiredService<IActivityRegistry>();
|
||||
var workflowProviders = Scope.ServiceProvider.GetRequiredService<IEnumerable<IWorkflowProvider>>();
|
||||
_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<MaterializedWorkflow> 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<IWorkflowDefinitionManager>();
|
||||
_workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionsReloader>();
|
||||
}
|
||||
|
||||
[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);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue