diff --git a/src/core/Elsa.Core/Helpers/Diff.cs b/src/core/Elsa.Core/Helpers/Diff.cs index 876932fc5..46c5a7377 100644 --- a/src/core/Elsa.Core/Helpers/Diff.cs +++ b/src/core/Elsa.Core/Helpers/Diff.cs @@ -19,11 +19,11 @@ public class Diff public static class Diff { - public static Diff For(ICollection firstSet, ICollection secondSet) + public static Diff For(ICollection firstSet, ICollection secondSet, IEqualityComparer? comparer = default) { - var removed = firstSet.Except(secondSet).ToList(); - var added = secondSet.Except(firstSet).ToList(); - var unchanged = firstSet.Intersect(secondSet).ToList(); + var removed = firstSet.Except(secondSet, comparer).ToList(); + var added = secondSet.Except(firstSet, comparer).ToList(); + var unchanged = firstSet.Intersect(secondSet, comparer).ToList(); return new Diff(added, removed, unchanged); } diff --git a/src/core/Elsa.Formatting/Formatters/JsonFormatter.cs b/src/core/Elsa.Formatting/Formatters/JsonFormatter.cs index 98687c3f1..16443aca8 100644 --- a/src/core/Elsa.Formatting/Formatters/JsonFormatter.cs +++ b/src/core/Elsa.Formatting/Formatters/JsonFormatter.cs @@ -1,4 +1,5 @@ using System.Text.Json; +using System.Text.Json.Serialization; using Elsa.Formatting.Contracts; namespace Elsa.Formatting.Formatters; @@ -8,12 +9,13 @@ public class JsonFormatter : IFormatter public ValueTask ToStringAsync(object body, CancellationToken cancellationToken) { var json = JsonSerializer.Serialize(body); - return ValueTask.FromResult(json); + return ValueTask.FromResult(json); } public ValueTask FromStringAsync(string data, Type? returnType, CancellationToken cancellationToken) { var options = new JsonSerializerOptions(); + options.Converters.Add(new JsonStringEnumConverter()); var value = returnType != null ? JsonSerializer.Deserialize(data, returnType, options)! : JsonSerializer.Deserialize(data, options)!; return ValueTask.FromResult(value); } diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Activities/MessageReceived.cs b/src/modules/Elsa.Modules.AzureServiceBus/Activities/MessageReceived.cs index 3cfa48432..a17673bef 100644 --- a/src/modules/Elsa.Modules.AzureServiceBus/Activities/MessageReceived.cs +++ b/src/modules/Elsa.Modules.AzureServiceBus/Activities/MessageReceived.cs @@ -7,9 +7,9 @@ using Elsa.Modules.AzureServiceBus.Models; namespace Elsa.Modules.AzureServiceBus.Activities; [Activity("Elsa.AzureServiceBus.MessageReceived", "Executes when a message is received from the configured queue or topic and subscription", "Azure Service Bus")] -public class MessageReceived : Trigger +public class MessageReceived : Trigger { - internal const string MessageReceivedInputKey = "ReceivedMessage"; + internal const string InputKey = "ReceivedMessage"; [JsonConstructor] public MessageReceived() @@ -48,11 +48,6 @@ public class MessageReceived : Trigger /// public Output? ReceivedMessage { get; set; } - /// - /// The parsed body of the received message. - /// - public Output? ReceivedMessageBody { get; set; } - /// /// The formatter to use to parse the message. /// @@ -60,21 +55,34 @@ public class MessageReceived : Trigger protected override object GetTriggerDatum(TriggerIndexingContext context) => GetBookmarkData(context.ExpressionExecutionContext); - protected override void Execute(ActivityExecutionContext context) + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { - var bookmarkData = GetBookmarkData(context.ExpressionExecutionContext); - context.CreateBookmark(bookmarkData, Resume); + // If we did not receive external input, it means we are just now encountering this activity. + if (!context.TryGetInput(InputKey, out var receivedMessage)) + { + // Create bookmarks for when we receive the expected HTTP request. + context.CreateBookmark(GetBookmarkData(context.ExpressionExecutionContext), Resume); + return; + } + + // Provide the received message as output. + await SetResultAsync(receivedMessage, context); } private async ValueTask Resume(ActivityExecutionContext context) { - var receivedMessage = (ReceivedServiceBusMessageModel)context.WorkflowExecutionContext.Input[MessageReceivedInputKey]!; + var receivedMessage = context.GetInput(InputKey); + await SetResultAsync(receivedMessage, context); + } + + private async Task SetResultAsync(ReceivedServiceBusMessageModel receivedMessage, ActivityExecutionContext context) + { var bodyAsString = new BinaryData(receivedMessage.Body).ToString(); var targetType = context.Get(ExpectedMessageType); var body = Formatter == null ? bodyAsString : await Formatter.FromStringAsync(bodyAsString, targetType, context.CancellationToken); context.Set(ReceivedMessage, receivedMessage); - context.Set(ReceivedMessageBody, body); + context.Set(Result, body); } private object GetBookmarkData(ExpressionExecutionContext context) diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IWorkerManager.cs b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IWorkerManager.cs index 16dfb6831..fcafa1fab 100644 --- a/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IWorkerManager.cs +++ b/src/modules/Elsa.Modules.AzureServiceBus/Contracts/IWorkerManager.cs @@ -29,6 +29,11 @@ public interface IWorkerManager /// Task StopWorkerAsync(string queueOrTopic, string? subscription, CancellationToken cancellationToken = default); + /// + /// Ensures that at least one worker exists for the specified queue/topic and subscription. + /// + Task EnsureWorkerAsync(string queueOrTopic, string? subscription, CancellationToken cancellationToken = default); + /// /// Removes the specified worker. /// diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Handlers/UpdateWorkers.cs b/src/modules/Elsa.Modules.AzureServiceBus/Handlers/UpdateWorkers.cs index f68f30c1a..116e8f0a5 100644 --- a/src/modules/Elsa.Modules.AzureServiceBus/Handlers/UpdateWorkers.cs +++ b/src/modules/Elsa.Modules.AzureServiceBus/Handlers/UpdateWorkers.cs @@ -30,9 +30,11 @@ public class UpdateWorkers : INotificationHandler, INot { var added = notification.IndexedWorkflowTriggers.AddedTriggers.Filter().Select(x => DeserializePayload(x.Data!)); var removed = notification.IndexedWorkflowTriggers.RemovedTriggers.Filter().Select(x => DeserializePayload(x.Data!)); + var unchanged = notification.IndexedWorkflowTriggers.UnchangedTriggers.Filter().Select(x => DeserializePayload(x.Data!)); await StopWorkersAsync(removed, cancellationToken); await StartWorkersAsync(added, cancellationToken); + await EnsureWorkersAsync(unchanged, cancellationToken); } /// @@ -65,4 +67,9 @@ public class UpdateWorkers : INotificationHandler, INot { foreach (var payload in payloads) await _workerManager.StopWorkerAsync(payload.QueueOrTopic, payload.Subscription, cancellationToken); } + + private async Task EnsureWorkersAsync(IEnumerable payloads, CancellationToken cancellationToken) + { + foreach (var payload in payloads) await _workerManager.EnsureWorkerAsync(payload.QueueOrTopic, payload.Subscription, cancellationToken); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Services/Worker.cs b/src/modules/Elsa.Modules.AzureServiceBus/Services/Worker.cs index b0b831813..46d7cac53 100644 --- a/src/modules/Elsa.Modules.AzureServiceBus/Services/Worker.cs +++ b/src/modules/Elsa.Modules.AzureServiceBus/Services/Worker.cs @@ -73,7 +73,8 @@ public class Worker : IAsyncDisposable var payload = new MessageReceivedTriggerPayload(QueueOrTopic, Subscription); var hash = _hasher.Hash(payload); var messageModel = CreateMessageModel(message); - var stimulus = Stimulus.Standard(BookmarkName, hash, new { ReceivedMessage = messageModel }); + var input = new Dictionary() { [MessageReceived.InputKey] = messageModel }; + var stimulus = Stimulus.Standard(BookmarkName, hash, input); var executionResults = (await _workflowService.ExecuteStimulusAsync(stimulus, cancellationToken)).ToList(); _logger.LogInformation("Triggered {WorkflowCount} workflows", executionResults.Count); diff --git a/src/modules/Elsa.Modules.AzureServiceBus/Services/WorkerManager.cs b/src/modules/Elsa.Modules.AzureServiceBus/Services/WorkerManager.cs index 942b5eb61..13830d21e 100644 --- a/src/modules/Elsa.Modules.AzureServiceBus/Services/WorkerManager.cs +++ b/src/modules/Elsa.Modules.AzureServiceBus/Services/WorkerManager.cs @@ -26,11 +26,7 @@ public class WorkerManager : IWorkerManager, IAsyncDisposable return; } - subscription ??= ""; - worker = ActivatorUtilities.CreateInstance(_serviceProvider, queueOrTopic, subscription!); - - _workers.Add(worker); - await worker.StartAsync(cancellationToken); + await CreateAndAddWorkerAsync(queueOrTopic, subscription, cancellationToken); } public async Task StopWorkerAsync(string queueOrTopic, string? subscription, CancellationToken cancellationToken = default) @@ -44,6 +40,13 @@ public class WorkerManager : IWorkerManager, IAsyncDisposable if (worker.RefCount == 0) await RemoveWorkerAsync(worker); } + + public async Task EnsureWorkerAsync(string queueOrTopic, string? subscription, CancellationToken cancellationToken = default) + { + var worker = FindWorkerFor(queueOrTopic, subscription); + if (worker != null) return; + await CreateAndAddWorkerAsync(queueOrTopic, subscription, cancellationToken); + } public async Task RemoveWorkerAsync(Worker worker) { @@ -55,4 +58,13 @@ public class WorkerManager : IWorkerManager, IAsyncDisposable { foreach (var worker in Workers) await worker.DisposeAsync(); } + + private async Task CreateAndAddWorkerAsync(string queueOrTopic, string? subscription, CancellationToken cancellationToken = default) + { + subscription ??= ""; + var worker = ActivatorUtilities.CreateInstance(_serviceProvider, queueOrTopic, subscription!); + + _workers.Add(worker); + await worker.StartAsync(cancellationToken); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Modules.Http/Middleware/HttpTriggerMiddleware.cs b/src/modules/Elsa.Modules.Http/Middleware/HttpTriggerMiddleware.cs index 863f20654..e01db1908 100644 --- a/src/modules/Elsa.Modules.Http/Middleware/HttpTriggerMiddleware.cs +++ b/src/modules/Elsa.Modules.Http/Middleware/HttpTriggerMiddleware.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Generic; using System.Linq; using System.Net.Mime; using System.Text.Json; @@ -33,7 +34,7 @@ public class HttpTriggerMiddleware var hash = _hasher.Hash(new HttpBookmarkData(path, method)); var activityTypeName = TypeNameHelper.GenerateTypeName(); var requestModel = new HttpRequestModel(new Uri(request.GetEncodedUrl())); - var input = new { HttpRequest = requestModel }; + var input = new Dictionary() { [HttpEndpoint.InputKey] = requestModel }; var stimulus = Stimulus.Standard(activityTypeName, hash, input); var executionResults = (await workflowService.ExecuteStimulusAsync(stimulus, abortToken)).ToList(); diff --git a/src/persistence/Elsa.Persistence.Abstractions/Comparers/WorkflowTriggerHashEqualityComparer.cs b/src/persistence/Elsa.Persistence.Abstractions/Comparers/WorkflowTriggerHashEqualityComparer.cs new file mode 100644 index 000000000..eadaa8441 --- /dev/null +++ b/src/persistence/Elsa.Persistence.Abstractions/Comparers/WorkflowTriggerHashEqualityComparer.cs @@ -0,0 +1,9 @@ +using Elsa.Persistence.Entities; + +namespace Elsa.Persistence.Comparers; + +public class WorkflowTriggerHashEqualityComparer : IEqualityComparer +{ + public bool Equals(WorkflowTrigger? x, WorkflowTrigger? y) => x?.Hash?.Equals(y?.Hash) ?? false; + public int GetHashCode(WorkflowTrigger obj) => obj.Hash?.GetHashCode() ?? "".GetHashCode(); +} \ No newline at end of file diff --git a/src/runtime/Elsa.Runtime/Models/IndexedWorkflowTriggers.cs b/src/runtime/Elsa.Runtime/Models/IndexedWorkflowTriggers.cs index 0fb57c034..ba005ee4c 100644 --- a/src/runtime/Elsa.Runtime/Models/IndexedWorkflowTriggers.cs +++ b/src/runtime/Elsa.Runtime/Models/IndexedWorkflowTriggers.cs @@ -3,4 +3,4 @@ using Elsa.Persistence.Entities; namespace Elsa.Runtime.Models; -public record IndexedWorkflowTriggers(Workflow Workflow, ICollection AddedTriggers, ICollection RemovedTriggers); \ No newline at end of file +public record IndexedWorkflowTriggers(Workflow Workflow, ICollection AddedTriggers, ICollection RemovedTriggers, ICollection UnchangedTriggers); \ No newline at end of file diff --git a/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs b/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs index 44a3f4daf..a09ff49f8 100644 --- a/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs +++ b/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs @@ -7,6 +7,7 @@ using Elsa.Helpers; using Elsa.Mediator.Contracts; using Elsa.Models; using Elsa.Persistence.Commands; +using Elsa.Persistence.Comparers; using Elsa.Persistence.Entities; using Elsa.Persistence.Requests; using Elsa.Runtime.Contracts; @@ -95,12 +96,12 @@ public class TriggerIndexer : ITriggerIndexer : new List(0); // Diff triggers. - var diff = Diff.For(currentTriggers, newTriggers); + var diff = Diff.For(currentTriggers, newTriggers, new WorkflowTriggerHashEqualityComparer()); // Replace triggers for the specified workflow. await _commandSender.ExecuteAsync(new ReplaceWorkflowTriggers(workflow, diff.Removed, diff.Added), cancellationToken); - var indexedWorkflow = new IndexedWorkflowTriggers(workflow, diff.Added, diff.Removed); + var indexedWorkflow = new IndexedWorkflowTriggers(workflow, diff.Added, diff.Removed, diff.Unchanged); // Publish event. await _eventPublisher.PublishAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken); diff --git a/src/runtime/Elsa.Runtime/WorkflowProviders/ConfigurationWorkflowProvider.cs b/src/runtime/Elsa.Runtime/WorkflowProviders/ConfigurationWorkflowProvider.cs index 13732885a..1e4517f6a 100644 --- a/src/runtime/Elsa.Runtime/WorkflowProviders/ConfigurationWorkflowProvider.cs +++ b/src/runtime/Elsa.Runtime/WorkflowProviders/ConfigurationWorkflowProvider.cs @@ -46,10 +46,11 @@ public class ConfigurationWorkflowProvider : IWorkflowProvider private Workflow BuildWorkflowDefinition(Func workflowFactory) { var builder = new WorkflowDefinitionBuilder(); - builder.WithDefinitionId(workflowFactory.GetType().Name); var definition = workflowFactory(_serviceProvider); + + builder.WithDefinitionId(definition.GetType().Name); definition.Build(builder); - + var workflow = builder.BuildWorkflow(); _identityGraphService.AssignIdentities(workflow); diff --git a/src/samples/aspnet/Elsa.Samples.Web1/Workflows/ReceiveMessageWorkflow.cs b/src/samples/aspnet/Elsa.Samples.Web1/Workflows/ReceiveMessageWorkflow.cs index 32fb1e090..7509cbb05 100644 --- a/src/samples/aspnet/Elsa.Samples.Web1/Workflows/ReceiveMessageWorkflow.cs +++ b/src/samples/aspnet/Elsa.Samples.Web1/Workflows/ReceiveMessageWorkflow.cs @@ -22,7 +22,7 @@ public class ReceiveMessageWorkflow : IWorkflow { CanStartWorkflow = true, QueueOrTopic = new Input("inbox"), - ReceivedMessageBody = new Output(receivedMessage) + Result = new Output(receivedMessage) }, new WriteLine(context => $"Message received: {receivedMessage.Get(context)}") }