Update AzureServiceBus module

This commit is contained in:
Sipke Schoorstra 2022-04-07 11:22:44 +02:00
parent 03b36393fe
commit e7c36ed7ff
13 changed files with 77 additions and 30 deletions

View file

@ -19,11 +19,11 @@ public class Diff<T>
public static class Diff
{
public static Diff<T> For<T>(ICollection<T> firstSet, ICollection<T> secondSet)
public static Diff<T> For<T>(ICollection<T> firstSet, ICollection<T> secondSet, IEqualityComparer<T>? 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<T>(added, removed, unchanged);
}

View file

@ -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<string> ToStringAsync(object body, CancellationToken cancellationToken)
{
var json = JsonSerializer.Serialize(body);
return ValueTask.FromResult<string>(json);
return ValueTask.FromResult(json);
}
public ValueTask<object> 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<object>(data, options)!;
return ValueTask.FromResult(value);
}

View file

@ -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<object>
{
internal const string MessageReceivedInputKey = "ReceivedMessage";
internal const string InputKey = "ReceivedMessage";
[JsonConstructor]
public MessageReceived()
@ -48,11 +48,6 @@ public class MessageReceived : Trigger
/// </summary>
public Output<ReceivedServiceBusMessageModel>? ReceivedMessage { get; set; }
/// <summary>
/// The parsed body of the received message.
/// </summary>
public Output<object>? ReceivedMessageBody { get; set; }
/// <summary>
/// The formatter to use to parse the message.
/// </summary>
@ -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<ReceivedServiceBusMessageModel>(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<ReceivedServiceBusMessageModel>(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)

View file

@ -29,6 +29,11 @@ public interface IWorkerManager
/// </summary>
Task StopWorkerAsync(string queueOrTopic, string? subscription, CancellationToken cancellationToken = default);
/// <summary>
/// Ensures that at least one worker exists for the specified queue/topic and subscription.
/// </summary>
Task EnsureWorkerAsync(string queueOrTopic, string? subscription, CancellationToken cancellationToken = default);
/// <summary>
/// Removes the specified worker.
/// </summary>

View file

@ -30,9 +30,11 @@ public class UpdateWorkers : INotificationHandler<WorkflowTriggersIndexed>, INot
{
var added = notification.IndexedWorkflowTriggers.AddedTriggers.Filter<MessageReceived>().Select(x => DeserializePayload(x.Data!));
var removed = notification.IndexedWorkflowTriggers.RemovedTriggers.Filter<MessageReceived>().Select(x => DeserializePayload(x.Data!));
var unchanged = notification.IndexedWorkflowTriggers.UnchangedTriggers.Filter<MessageReceived>().Select(x => DeserializePayload(x.Data!));
await StopWorkersAsync(removed, cancellationToken);
await StartWorkersAsync(added, cancellationToken);
await EnsureWorkersAsync(unchanged, cancellationToken);
}
/// <summary>
@ -65,4 +67,9 @@ public class UpdateWorkers : INotificationHandler<WorkflowTriggersIndexed>, INot
{
foreach (var payload in payloads) await _workerManager.StopWorkerAsync(payload.QueueOrTopic, payload.Subscription, cancellationToken);
}
private async Task EnsureWorkersAsync(IEnumerable<MessageReceivedTriggerPayload> payloads, CancellationToken cancellationToken)
{
foreach (var payload in payloads) await _workerManager.EnsureWorkerAsync(payload.QueueOrTopic, payload.Subscription, cancellationToken);
}
}

View file

@ -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<string, object>() { [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);

View file

@ -26,11 +26,7 @@ public class WorkerManager : IWorkerManager, IAsyncDisposable
return;
}
subscription ??= "";
worker = ActivatorUtilities.CreateInstance<Worker>(_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<Worker>(_serviceProvider, queueOrTopic, subscription!);
_workers.Add(worker);
await worker.StartAsync(cancellationToken);
}
}

View file

@ -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<HttpEndpoint>();
var requestModel = new HttpRequestModel(new Uri(request.GetEncodedUrl()));
var input = new { HttpRequest = requestModel };
var input = new Dictionary<string, object>() { [HttpEndpoint.InputKey] = requestModel };
var stimulus = Stimulus.Standard(activityTypeName, hash, input);
var executionResults = (await workflowService.ExecuteStimulusAsync(stimulus, abortToken)).ToList();

View file

@ -0,0 +1,9 @@
using Elsa.Persistence.Entities;
namespace Elsa.Persistence.Comparers;
public class WorkflowTriggerHashEqualityComparer : IEqualityComparer<WorkflowTrigger>
{
public bool Equals(WorkflowTrigger? x, WorkflowTrigger? y) => x?.Hash?.Equals(y?.Hash) ?? false;
public int GetHashCode(WorkflowTrigger obj) => obj.Hash?.GetHashCode() ?? "".GetHashCode();
}

View file

@ -3,4 +3,4 @@ using Elsa.Persistence.Entities;
namespace Elsa.Runtime.Models;
public record IndexedWorkflowTriggers(Workflow Workflow, ICollection<WorkflowTrigger> AddedTriggers, ICollection<WorkflowTrigger> RemovedTriggers);
public record IndexedWorkflowTriggers(Workflow Workflow, ICollection<WorkflowTrigger> AddedTriggers, ICollection<WorkflowTrigger> RemovedTriggers, ICollection<WorkflowTrigger> UnchangedTriggers);

View file

@ -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<WorkflowTrigger>(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);

View file

@ -46,10 +46,11 @@ public class ConfigurationWorkflowProvider : IWorkflowProvider
private Workflow BuildWorkflowDefinition(Func<IServiceProvider, IWorkflow> 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);

View file

@ -22,7 +22,7 @@ public class ReceiveMessageWorkflow : IWorkflow
{
CanStartWorkflow = true,
QueueOrTopic = new Input<string>("inbox"),
ReceivedMessageBody = new Output<object>(receivedMessage)
Result = new Output<object?>(receivedMessage)
},
new WriteLine(context => $"Message received: {receivedMessage.Get(context)}")
}