diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index 05a7ed157..440d9c904 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -5,11 +5,16 @@ on: branches: - 'main' - 'feature/*' + - 'feat/*' - 'issue/*' - 'bug/*' - 'enhancement/*' + - 'enh/*' - 'patch/*' - 'fix/*' + - 'perf/*' + - 'hotfix/*' + - 'chore/*' release: types: [ prereleased, published ] env: diff --git a/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs b/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs index 53cc92daa..8f1845ac6 100644 --- a/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs +++ b/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs @@ -6,6 +6,7 @@ using Elsa.Workflows; using Elsa.Workflows.Attributes; using Elsa.Workflows.UIHints; using Elsa.Workflows.Models; +using Microsoft.Extensions.Logging; using HttpHeaders = Elsa.Http.Models.HttpHeaders; namespace Elsa.Http; @@ -119,6 +120,7 @@ public abstract class SendHttpRequestBase : Activity private async Task TrySendAsync(ActivityExecutionContext context) { var request = PrepareRequest(context); + var logger = (ILogger)context.GetRequiredService(typeof(ILogger<>).MakeGenericType(GetType())); var httpClientFactory = context.GetRequiredService(); var httpClient = httpClientFactory.CreateClient(nameof(SendHttpRequestBase)); var cancellationToken = context.CancellationToken; @@ -139,12 +141,14 @@ public abstract class SendHttpRequestBase : Activity } catch (HttpRequestException e) { + logger.LogWarning(e, "An error occurred while sending an HTTP request"); context.AddExecutionLogEntry("Error", e.Message, payload: new { StackTrace = e.StackTrace }); context.JournalData.Add("Error", e.Message); await HandleRequestExceptionAsync(context, e); } catch (TaskCanceledException e) { + logger.LogWarning(e, "An error occurred while sending an HTTP request"); context.AddExecutionLogEntry("Error", e.Message, payload: new { StackTrace = e.StackTrace }); context.JournalData.Add("Cancelled", true); await HandleTaskCanceledExceptionAsync(context, e); diff --git a/src/modules/Elsa.MassTransit/Messages/DispatchResumeWorkflows.cs b/src/modules/Elsa.MassTransit/Messages/DispatchResumeWorkflows.cs index 2754eee1e..9d059d277 100644 --- a/src/modules/Elsa.MassTransit/Messages/DispatchResumeWorkflows.cs +++ b/src/modules/Elsa.MassTransit/Messages/DispatchResumeWorkflows.cs @@ -3,6 +3,7 @@ using Elsa.Workflows.Serialization.Converters; namespace Elsa.MassTransit.Messages; +[Obsolete("This message is no longer used and will be removed in a future version.")] public class DispatchResumeWorkflows(string activityTypeName, object bookmarkPayload) { public string ActivityTypeName { get; init; } = activityTypeName; diff --git a/src/modules/Elsa.MassTransit/Messages/DispatchTriggerWorkflowsRequest.cs b/src/modules/Elsa.MassTransit/Messages/DispatchTriggerWorkflowsRequest.cs index 9f57d3b69..36a140b13 100644 --- a/src/modules/Elsa.MassTransit/Messages/DispatchTriggerWorkflowsRequest.cs +++ b/src/modules/Elsa.MassTransit/Messages/DispatchTriggerWorkflowsRequest.cs @@ -3,6 +3,7 @@ using Elsa.Workflows.Serialization.Converters; namespace Elsa.MassTransit.Messages; +[Obsolete("This is no longer used and will be removed in a future version.")] public class DispatchTriggerWorkflows(string activityTypeName, object bookmarkPayload) { public string ActivityTypeName { get; init; } = activityTypeName; diff --git a/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs b/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs index c5b23ce35..4da57240d 100644 --- a/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs +++ b/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs @@ -1,12 +1,17 @@ +using Elsa.Extensions; using Elsa.MassTransit.Contracts; using Elsa.MassTransit.Messages; +using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Requests; using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Entities; +using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Requests; using Elsa.Workflows.Runtime.Responses; using MassTransit; +using Microsoft.Extensions.Logging; namespace Elsa.MassTransit.Services; @@ -17,18 +22,16 @@ public class MassTransitWorkflowDispatcher( IBus bus, IEndpointChannelFormatter endpointChannelFormatter, IWorkflowDefinitionService workflowDefinitionService, - IWorkflowInstanceManager workflowInstanceManager) + IWorkflowInstanceManager workflowInstanceManager, + IBookmarkHasher bookmarkHasher, + ITriggerStore triggerStore, + IBookmarkStore bookmarkStore, + ILogger logger) : IWorkflowDispatcher { /// public async Task DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) { - // When a request is received to execute a workflow, the initial step taken by our system is the creation of the workflow instance. - // The input parameters for the particular instance are immediately persisted in the database. This step is crucial for a couple of reasons: - // 1. Size constraint: It helps us prevent scenarios where the message size exceeds limits. Large messages cause the system to fail at the sending stage. - // 2. Performance: A smaller message size means less information needs to be processed and transferred, optimizing speed and efficiency. - - // To create the instance, we need to find the workflow definition first. var workflow = await workflowDefinitionService.FindWorkflowAsync(request.DefinitionId, request.VersionOptions, cancellationToken); if (workflow == null) @@ -44,13 +47,7 @@ public class MassTransitWorkflowDispatcher( CorrelationId = request.CorrelationId }; - // The workflow instance is created and persisted in the database. - var workflowInstance = await workflowInstanceManager.CreateWorkflowInstanceAsync(createWorkflowInstanceRequest, cancellationToken); - - // The workflow instance is then dispatched for execution. - var sendEndpoint = await GetSendEndpointAsync(options); - var message = DispatchWorkflowDefinition.DispatchExistingWorkflowInstance(workflowInstance.Id, request.TriggerActivityId); - await sendEndpoint.Send(message, cancellationToken); + await DispatchWorkflowAsync(createWorkflowInstanceRequest, request.TriggerActivityId, options, cancellationToken); return DispatchWorkflowResponse.Success(); } @@ -74,31 +71,120 @@ public class MassTransitWorkflowDispatcher( /// public async Task DispatchAsync(DispatchTriggerWorkflowsRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) { - var sendEndpoint = await GetSendEndpointAsync(options); - await sendEndpoint.Send(new DispatchTriggerWorkflows(request.ActivityTypeName, request.BookmarkPayload) - { - CorrelationId = request.CorrelationId, - WorkflowInstanceId = request.WorkflowInstanceId, - ActivityInstanceId = request.ActivityInstanceId, - Input = request.Input - }, cancellationToken); + await DispatchTriggersAsync(request, options, cancellationToken); + await DispatchBookmarksAsync(request, options, cancellationToken); return DispatchWorkflowResponse.Success(); } /// public async Task DispatchAsync(DispatchResumeWorkflowsRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) { - var sendEndpoint = await GetSendEndpointAsync(options); - await sendEndpoint.Send(new DispatchResumeWorkflows(request.ActivityTypeName, request.BookmarkPayload) + var hash = bookmarkHasher.Hash(request.ActivityTypeName, request.BookmarkPayload, request.ActivityInstanceId); + var correlationId = request.CorrelationId; + var workflowInstanceId = request.WorkflowInstanceId; + var activityInstanceId = request.ActivityInstanceId; + var filter = new BookmarkFilter { - CorrelationId = request.CorrelationId, - WorkflowInstanceId = request.WorkflowInstanceId, - ActivityInstanceId = request.ActivityInstanceId, - Input = request.Input - }, cancellationToken); + Hash = hash, + CorrelationId = correlationId, + WorkflowInstanceId = workflowInstanceId, + ActivityInstanceId = activityInstanceId + }; + var bookmarks = await bookmarkStore.FindManyAsync(filter, cancellationToken); + await DispatchBookmarksAsync(bookmarks, request.Input, null, options, cancellationToken); return DispatchWorkflowResponse.Success(); } + private async Task DispatchTriggersAsync(DispatchTriggerWorkflowsRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) + { + var triggerHash = bookmarkHasher.Hash(request.ActivityTypeName, request.BookmarkPayload); + var triggerFilter = new TriggerFilter + { + Hash = triggerHash + }; + var triggers = (await triggerStore.FindManyAsync(triggerFilter, cancellationToken)).ToList(); + + foreach (var trigger in triggers) + { + var workflow = await workflowDefinitionService.FindWorkflowAsync(trigger.WorkflowDefinitionVersionId, cancellationToken); + + if (workflow == null) + { + logger.LogWarning("Workflow definition with ID '{WorkflowDefinitionId}' not found", trigger.WorkflowDefinitionVersionId); + continue; + } + + var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest + { + Workflow = workflow, + WorkflowInstanceId = request.WorkflowInstanceId, + Input = request.Input, + Properties = request.Properties, + CorrelationId = request.CorrelationId + }; + + await DispatchWorkflowAsync(createWorkflowInstanceRequest, trigger.ActivityId, options, cancellationToken); + } + } + + private async Task DispatchWorkflowAsync(CreateWorkflowInstanceRequest createWorkflowInstanceRequest, string? triggerActivityId, DispatchWorkflowOptions? options, CancellationToken cancellationToken) + { + var workflowInstance = await workflowInstanceManager.CreateWorkflowInstanceAsync(createWorkflowInstanceRequest, cancellationToken); + var sendEndpoint = await GetSendEndpointAsync(options); + var message = DispatchWorkflowDefinition.DispatchExistingWorkflowInstance(workflowInstance.Id, triggerActivityId); + await sendEndpoint.Send(message, cancellationToken); + } + + private async Task DispatchBookmarksAsync(DispatchTriggerWorkflowsRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) + { + var correlationId = request.CorrelationId; + var workflowInstanceId = request.WorkflowInstanceId; + var activityInstanceId = request.ActivityInstanceId; + var bookmarkHash = bookmarkHasher.Hash(request.ActivityTypeName, request.BookmarkPayload, activityInstanceId); + + var filter = new BookmarkFilter + { + Hash = bookmarkHash, + CorrelationId = correlationId, + WorkflowInstanceId = workflowInstanceId, + ActivityInstanceId = activityInstanceId + }; + var bookmarks = (await bookmarkStore.FindManyAsync(filter, cancellationToken)).ToList(); + + await DispatchBookmarksAsync(bookmarks, request.Input, request.Properties, options, cancellationToken); + } + + private async Task DispatchBookmarksAsync(IEnumerable bookmarks, IDictionary? input, IDictionary? properties, DispatchWorkflowOptions? options, CancellationToken cancellationToken) + { + foreach (var bookmark in bookmarks) + { + var workflowInstanceId = bookmark.WorkflowInstanceId; + + if (input != null || properties != null) + { + var workflowInstance = await workflowInstanceManager.FindByIdAsync(workflowInstanceId, cancellationToken); + + if (workflowInstance == null) + { + logger.LogWarning("Workflow instance with ID '{WorkflowInstanceId}' not found", workflowInstanceId); + continue; + } + + if (input != null) workflowInstance.WorkflowState.Input.Merge(input); + if (properties != null) workflowInstance.WorkflowState.Properties.Merge(properties); + + await workflowInstanceManager.SaveAsync(workflowInstance, cancellationToken); + } + + var dispatchInstanceRequest = new DispatchWorkflowInstanceRequest(workflowInstanceId) + { + BookmarkId = bookmark.BookmarkId, + CorrelationId = bookmark.CorrelationId + }; + await DispatchAsync(dispatchInstanceRequest, options, cancellationToken); + } + } + private async Task GetSendEndpointAsync(DispatchWorkflowOptions? options = default) { var endpointName = endpointChannelFormatter.FormatEndpointName(options?.Channel); diff --git a/src/modules/Elsa.Workflows.Core/Serialization/Converters/PolymorphicObjectConverter.cs b/src/modules/Elsa.Workflows.Core/Serialization/Converters/PolymorphicObjectConverter.cs index 2efd594d0..c5fd590c7 100644 --- a/src/modules/Elsa.Workflows.Core/Serialization/Converters/PolymorphicObjectConverter.cs +++ b/src/modules/Elsa.Workflows.Core/Serialization/Converters/PolymorphicObjectConverter.cs @@ -1,5 +1,7 @@ using System.Collections; using System.Dynamic; +using System.Reflection; +using System.Runtime; using System.Text.Json; using System.Text.Json.Nodes; using System.Text.Json.Serialization; @@ -47,7 +49,7 @@ public class PolymorphicObjectConverter : JsonConverter { return JsonSerializer.Deserialize(ref reader, targetType, newOptions)!; } - catch (NotSupportedException e) + catch (Exception e) when (e is NotSupportedException or TargetException) { return default!; } diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs index 7e9d1c427..79cf91274 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs @@ -85,8 +85,9 @@ public class WorkflowDefinitionActivity : Composite, IInitializable private void CopyInputOutputToVariables(ActivityExecutionContext context) { var serviceProvider = context.GetRequiredService(); + var activityDescriptor = FindActivityDescriptor(serviceProvider); - DeclareInputAsVariables(serviceProvider, (descriptor, variable) => + DeclareInputAsVariables(activityDescriptor, (descriptor, variable) => { var inputName = descriptor.Name; var input = SyntheticProperties.TryGetValue(inputName, out var inputValue) ? (Input?)inputValue : default; @@ -96,14 +97,11 @@ public class WorkflowDefinitionActivity : Composite, IInitializable variable.Set(context, evaluatedExpression); }); - DeclareOutputAsVariables(serviceProvider, (descriptor, variable) => context.ExpressionExecutionContext.Memory.Declare(variable)); + DeclareOutputAsVariables(activityDescriptor, (descriptor, variable) => context.ExpressionExecutionContext.Memory.Declare(variable)); } - private void DeclareInputAsVariables(IServiceProvider serviceProvider, Action configureVariable) + private void DeclareInputAsVariables(ActivityDescriptor activityDescriptor, Action configureVariable) { - var activityRegistry = serviceProvider.GetRequiredService(); - var activityDescriptor = activityRegistry.Find(Type, Version)!; - foreach (var inputDescriptor in activityDescriptor.Inputs) { var inputName = inputDescriptor.Name; @@ -119,11 +117,8 @@ public class WorkflowDefinitionActivity : Composite, IInitializable } } - private void DeclareOutputAsVariables(IServiceProvider serviceProvider, Action configureVariable) + private void DeclareOutputAsVariables(ActivityDescriptor activityDescriptor, Action configureVariable) { - var activityRegistry = serviceProvider.GetRequiredService(); - var activityDescriptor = activityRegistry.Find(Type, Version)!; - foreach (var outputDescriptor in activityDescriptor.Outputs) { var outputName = outputDescriptor.Name; @@ -156,6 +151,12 @@ public class WorkflowDefinitionActivity : Composite, IInitializable return workflow; } + + private ActivityDescriptor FindActivityDescriptor(IServiceProvider serviceProvider) + { + var activityRegistry = serviceProvider.GetRequiredService(); + return activityRegistry.Find(Type, Version) ?? activityRegistry.Find(Type) ?? throw new Exception($"Could not find activity descriptor for {Type}."); + } async ValueTask IInitializable.InitializeAsync(InitializationContext context) { @@ -166,9 +167,11 @@ public class WorkflowDefinitionActivity : Composite, IInitializable if (workflow == null) throw new Exception($"Could not find workflow definition with ID {WorkflowDefinitionId}."); + var activityDescriptor = FindActivityDescriptor(serviceProvider); + // Declare input and output variables. - DeclareInputAsVariables(serviceProvider, (_, variable) => Variables.Declare(variable)); - DeclareOutputAsVariables(serviceProvider, (_, variable) => Variables.Declare(variable)); + DeclareInputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable)); + DeclareOutputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable)); // Set the root activity. Root = workflow; diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceManager.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceManager.cs index 627892a0b..27c91a79b 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceManager.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowInstanceManager.cs @@ -10,6 +10,11 @@ namespace Elsa.Workflows.Management.Contracts; /// public interface IWorkflowInstanceManager { + /// + /// Retrieves the workflow instance with the specified ID. + /// + Task FindByIdAsync(string id, CancellationToken cancellationToken = default); + /// /// Saves the specified workflow instance. /// diff --git a/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs b/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs index dcf65e950..4943f0fef 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs @@ -13,31 +13,31 @@ namespace Elsa.Workflows.Management.Handlers; /// [UsedImplicitly] internal class EvictWorkflowDefinitionServiceCache(IWorkflowDefinitionCacheManager workflowDefinitionCacheManager) : - INotificationHandler, - INotificationHandler, - INotificationHandler, - INotificationHandler + INotificationHandler, + INotificationHandler, + INotificationHandler, + INotificationHandler { /// - public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) + public async Task HandleAsync(WorkflowDefinitionPublishing notification, CancellationToken cancellationToken) { await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(notification.WorkflowDefinition.DefinitionId, cancellationToken); } /// - public async Task HandleAsync(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) + public async Task HandleAsync(WorkflowDefinitionRetracting notification, CancellationToken cancellationToken) { await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(notification.WorkflowDefinition.DefinitionId, cancellationToken); } /// - public async Task HandleAsync(WorkflowDefinitionDeleted notification, CancellationToken cancellationToken) + public async Task HandleAsync(WorkflowDefinitionDeleting notification, CancellationToken cancellationToken) { await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(notification.DefinitionId, cancellationToken); } /// - public async Task HandleAsync(WorkflowDefinitionsDeleted notification, CancellationToken cancellationToken) + public async Task HandleAsync(WorkflowDefinitionsDeleting notification, CancellationToken cancellationToken) { foreach (var definitionId in notification.DefinitionIds) await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(definitionId, cancellationToken); diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs index 8a0b625ff..3020faf67 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs @@ -51,21 +51,14 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService /// public async Task FindWorkflowDefinitionAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter - { - DefinitionId = definitionId, - VersionOptions = versionOptions - }; + var filter = new WorkflowDefinitionFilter { DefinitionId = definitionId, VersionOptions = versionOptions }; return await _workflowDefinitionStore.FindAsync(filter, cancellationToken); } /// public async Task FindWorkflowDefinitionAsync(string definitionVersionId, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter - { - Id = definitionVersionId - }; + var filter = new WorkflowDefinitionFilter { Id = definitionVersionId }; return await _workflowDefinitionStore.FindAsync(filter, cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs index d3112f82e..9a98f3ddb 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs @@ -1,4 +1,4 @@ -using Elsa.Common.Contracts; +using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; @@ -8,6 +8,7 @@ using Elsa.Workflows.Management.Mappers; using Elsa.Workflows.Management.Notifications; using Elsa.Workflows.Management.Requests; using Elsa.Workflows.State; +using Exception = System.Exception; namespace Elsa.Workflows.Management.Services; @@ -21,6 +22,12 @@ public class WorkflowInstanceManager( IWorkflowStateSerializer workflowStateSerializer) : IWorkflowInstanceManager { + /// + public async Task FindByIdAsync(string id, CancellationToken cancellationToken = default) + { + return await store.FindAsync(id, cancellationToken); + } + /// public async Task SaveAsync(WorkflowInstance workflowInstance, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs index 62f86a3c7..d9326097f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs @@ -13,6 +13,13 @@ public interface IWorkflowDefinitionStorePopulator /// /// The cancellation token. 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); /// /// Adds a workflow definition to the store. @@ -20,4 +27,12 @@ public interface IWorkflowDefinitionStorePopulator /// A materialized workflow. /// An optional cancellation token. Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default); + + /// + /// Adds a workflow definition to the store. + /// + /// A materialized workflow. + /// /// Whether to index triggers. + /// An optional cancellation token. + Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs index 939511d8a..642e4155b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs @@ -59,7 +59,7 @@ public class PersistWorkflowExecutionLogMiddleware : WorkflowExecutionMiddleware }).ToList(); await _workflowExecutionLogStore.AddManyAsync(entries, context.CancellationTokens.SystemCancellationToken); - + // Publish notification. await _notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationTokens.SystemCancellationToken); } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultRegistriesPopulator.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultRegistriesPopulator.cs index 7018bac79..e1c87dace 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultRegistriesPopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultRegistriesPopulator.cs @@ -27,7 +27,7 @@ public class DefaultRegistriesPopulator : IRegistriesPopulator await _activityRegistryPopulator.PopulateRegistryAsync(cancellationToken); // Stage 2: Populate the workflow definition store. - await _workflowDefinitionStorePopulator.PopulateStoreAsync(cancellationToken); + await _workflowDefinitionStorePopulator.PopulateStoreAsync(false, cancellationToken); // Stage 3: Re-populate the activity registry. // After the workflow definition store has been populated, we need to re-populate the activity registry to make sure that the activity descriptors are up-to-date. @@ -35,6 +35,6 @@ public class DefaultRegistriesPopulator : IRegistriesPopulator // Stage 4. Re-update the workflow definition store with the current set of activities. // Finally, we need to re-populate the workflow definition store to make sure that the workflow definitions are up-to-date. - await _workflowDefinitionStorePopulator.PopulateStoreAsync(cancellationToken); + await _workflowDefinitionStorePopulator.PopulateStoreAsync(true, cancellationToken); } } \ 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 3dc97a6be..ecd15938e 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -49,23 +49,37 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } /// - public async Task PopulateStoreAsync(CancellationToken cancellationToken = default) + public Task PopulateStoreAsync(CancellationToken cancellationToken = default) + { + return PopulateStoreAsync(true, cancellationToken); + } + + /// + public async Task PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default) { var providers = _workflowDefinitionProviders(); foreach (var provider in providers) { var results = await provider.GetWorkflowsAsync(cancellationToken).AsTask().ToList(); - foreach (var result in results) await AddAsync(result, cancellationToken); + foreach (var result in results) await AddAsync(result, indexTriggers, cancellationToken); } } /// - public async 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) { await AssignIdentities(materializedWorkflow.Workflow, cancellationToken); await AddOrUpdateAsync(materializedWorkflow, cancellationToken); - await IndexTriggersAsync(materializedWorkflow, cancellationToken); + + if (indexTriggers) + await IndexTriggersAsync(materializedWorkflow, cancellationToken); } private async Task AssignIdentities(Workflow workflow, CancellationToken cancellationToken) @@ -145,33 +159,32 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP workflowDefinition.ProviderName = materializedWorkflow.ProviderName; workflowDefinition.MaterializerContext = materializerContextJson; workflowDefinition.MaterializerName = materializedWorkflow.MaterializerName; - - if (existingDefinitionVersion is null - && workflowDefinitionsToSave.Any(w => w.Id == workflowDefinition.Id)) + + if (existingDefinitionVersion is null && workflowDefinitionsToSave.Any(w => w.Id == workflowDefinition.Id)) { - _logger.LogError("Trying to create a new workflow with existing id {workflowId}", workflowDefinition.Id); + _logger.LogInformation("Workflow with ID {WorkflowId} already exists", workflowDefinition.Id); return; } - + workflowDefinitionsToSave.Add(workflowDefinition); - + var duplicates = workflowDefinitionsToSave.GroupBy(wd => wd.Id) .Where(g => g.Count() > 1) .Select(g => g.Key) .ToList(); - + if (duplicates.Any()) { throw new Exception($"Unable to update WorkflowDefinition with ids {string.Join(',', duplicates)} multiple times."); } - + await _workflowDefinitionStore.SaveManyAsync(workflowDefinitionsToSave, cancellationToken); return; async Task UpdateIsLatest() { // Always try to update the IsLatest property based on the VersionNumber - + // Reset current latest definitions. var filter = new WorkflowDefinitionFilter { diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowInbox.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowInbox.cs index fb1017356..f663c9228 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowInbox.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowInbox.cs @@ -127,15 +127,15 @@ public class DefaultWorkflowInbox : IWorkflowInbox return new DeliverWorkflowInboxMessageResult(results.TriggeredWorkflows); } - - await _workflowDispatcher.DispatchAsync(new DispatchTriggerWorkflowsRequest(activityTypeName, bookmarkPayload) + + var dispatchRequest = new DispatchTriggerWorkflowsRequest(activityTypeName, bookmarkPayload) { CorrelationId = correlationId, WorkflowInstanceId = workflowInstanceId, ActivityInstanceId = activityInstanceId, Input = input - }, cancellationToken: cancellationToken); - + }; + await _workflowDispatcher.DispatchAsync(dispatchRequest, cancellationToken); return new DeliverWorkflowInboxMessageResult(new List()); }