revert MessageBodyExtensions and Add AzureServiceBusSendActivity Base Class for Metadata.

This commit is contained in:
jdevillard 2021-06-03 21:24:49 +02:00 committed by Sipke Schoorstra
parent f28c9a8be5
commit af94019f15
4 changed files with 112 additions and 65 deletions

View file

@ -0,0 +1,100 @@
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Design;
using Elsa.Expressions;
using Elsa.Serialization;
using Elsa.Services;
using Elsa.Services.Models;
using Microsoft.Azure.ServiceBus;
using Microsoft.Azure.ServiceBus.Core;
using System;
using System.Collections.Generic;
using System.Text;
using System.Threading.Tasks;
namespace Elsa.Activities.AzureServiceBus
{
public abstract class AzureServiceBusSendActivity : Activity
{
protected ISenderClient Sender;
protected IContentSerializer _serializer;
public AzureServiceBusSendActivity(IContentSerializer serializer)
{
_serializer = serializer;
}
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid, SyntaxNames.Json })]
public object Message { get; set; } = default!;
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? CorrelationId { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string ContentType { get; set; } = default!;
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? Label { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? To { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string MessageId { get; set; } = default!;
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? PartitionKey { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? ViaPartitionKey { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? ReplyTo { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? SessionId { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public DateTime ExpiresAtUtc { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public TimeSpan TimeToLive { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string? ReplyToSessionId { get; set; }
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public DateTime ScheduledEnqueueTimeUtc { get; set; }
[ActivityInput(DefaultSyntax = SyntaxNames.Json, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid,SyntaxNames.Json },UIHint = ActivityInputUIHints.MultiLine)]
public IDictionary<string,Object> UserProperties { get; set; } = new Dictionary<string, object>();
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context)
{
var message = CreateMessage(Message);
if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId))
message.CorrelationId = context.WorkflowExecutionContext.CorrelationId;
await Sender.SendAsync(message);
return Done();
}
protected Message CreateMessage(Object input)
{
var message = Extensions.MessageBodyExtensions.CreateMessage(_serializer, input);
message.CorrelationId = CorrelationId;
message.ContentType = ContentType;
message.Label = Label;
message.To = To;
message.PartitionKey = PartitionKey;
message.ViaPartitionKey = ViaPartitionKey;
message.ReplyTo = ReplyTo;
message.SessionId = SessionId;
message.ReplyToSessionId = ReplyToSessionId;
if (MessageId != null)
message.MessageId = MessageId;
if (TimeToLive != null && TimeToLive > TimeSpan.Zero)
message.TimeToLive = TimeToLive;
if (ScheduledEnqueueTimeUtc != null)
message.ScheduledEnqueueTimeUtc = ScheduledEnqueueTimeUtc;
if (UserProperties != null)
foreach (var props in UserProperties)
message.UserProperties.Add(props.Key, props.Value);
return message;
}
}
}

View file

@ -1,4 +1,4 @@
using System.Threading.Tasks;
using System.Threading.Tasks;
using Elsa.Activities.AzureServiceBus.Services;
using Elsa.ActivityResults;
using Elsa.Attributes;
@ -10,33 +10,23 @@ using Elsa.Services.Models;
namespace Elsa.Activities.AzureServiceBus
{
[Action(Category = "Azure Service Bus", DisplayName = "Send Service Bus Message", Description = "Sends a message to the specified queue", Outcomes = new[] { OutcomeNames.Done })]
public class SendAzureServiceBusQueueMessage : Activity
public class SendAzureServiceBusQueueMessage : AzureServiceBusSendActivity
{
private readonly IQueueMessageSenderFactory _queueMessageSenderFactory;
private readonly IContentSerializer _serializer;
public SendAzureServiceBusQueueMessage(IQueueMessageSenderFactory queueMessageSenderFactory, IContentSerializer serializer)
:base(serializer)
{
_queueMessageSenderFactory = queueMessageSenderFactory;
_serializer = serializer;
}
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
public string QueueName { get; set; } = default!;
[ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid, SyntaxNames.Json })]
public object Message { get; set; } = default!;
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context)
{
var sender = await _queueMessageSenderFactory.GetSenderAsync(QueueName, context.CancellationToken);
var message = Extensions.MessageBodyExtensions.CreateMessage(_serializer, Message);
if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId))
message.CorrelationId = context.WorkflowExecutionContext.CorrelationId;
await sender.SendAsync(message);
return Done();
Sender = await _queueMessageSenderFactory.GetSenderAsync(QueueName, context.CancellationToken);
return await base.OnExecuteAsync(context);
}
}
}

View file

@ -9,31 +9,23 @@ using Elsa.Services.Models;
namespace Elsa.Activities.AzureServiceBus
{
[Action(Category = "Azure Service Bus", DisplayName = "Send Service Bus Topic Message", Description = "Sends a message to the specified topic", Outcomes = new[] { OutcomeNames.Done })]
public class SendAzureServiceBusTopicMessage : Activity
public class SendAzureServiceBusTopicMessage : AzureServiceBusSendActivity
{
private readonly ITopicMessageSenderFactory _messageSenderFactory;
private readonly IContentSerializer _serializer;
public SendAzureServiceBusTopicMessage(ITopicMessageSenderFactory messageSenderFactory, IContentSerializer serializer)
:base(serializer)
{
_messageSenderFactory = messageSenderFactory;
_serializer = serializer;
_messageSenderFactory = messageSenderFactory;
}
[ActivityInput] public string TopicName { get; set; } = default!;
[ActivityInput] public object Message { get; set; } = default!;
[ActivityInput] public string TopicName { get; set; } = default!;
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context)
{
var sender = await _messageSenderFactory.GetTopicSenderAsync(TopicName, context.CancellationToken);
Sender = await _messageSenderFactory.GetTopicSenderAsync(TopicName, context.CancellationToken);
var message = Extensions.MessageBodyExtensions.CreateMessage(_serializer,Message);
if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId))
message.CorrelationId = context.WorkflowExecutionContext.CorrelationId;
await sender.SendAsync(message);
return Done();
return await base.OnExecuteAsync(context);
}
}
}

View file

@ -24,12 +24,7 @@ namespace Elsa.Activities.AzureServiceBus.Extensions
{
byte[] messageBytes;
if(message is MessageModel)
{
var messageModel = (MessageModel)message;
return CreateMessageFromMessageModel(messageModel);
}
else if (message is string s)
if (message is string s)
messageBytes = Encoding.UTF8.GetBytes(s);
else
{
@ -39,35 +34,5 @@ namespace Elsa.Activities.AzureServiceBus.Extensions
return new Message(messageBytes);
}
private static Message CreateMessageFromMessageModel(MessageModel message)
{
var returnMessage = new Message(message.Body)
{
CorrelationId = message.CorrelationId,
ContentType = message.ContentType,
Label = message.Label,
To = message.To,
PartitionKey = message.PartitionKey,
ViaPartitionKey = message.ViaPartitionKey,
ReplyTo = message.ReplyTo,
SessionId = message.SessionId,
ReplyToSessionId = message.ReplyToSessionId,
};
if(message.MessageId != null)
returnMessage.MessageId = message.MessageId;
if(message.TimeToLive != null && message.TimeToLive > TimeSpan.Zero )
returnMessage.TimeToLive = message.TimeToLive;
if(message.ScheduledEnqueueTimeUtc != null)
returnMessage.ScheduledEnqueueTimeUtc = message.ScheduledEnqueueTimeUtc;
if (message.UserProperties != null)
foreach (var props in message.UserProperties)
returnMessage.UserProperties.Add(props.Key, props.Value);
return returnMessage;
}
}
}