2020-11-22 21:16:28 +00:00
using System.Text ;
using System.Threading.Tasks ;
2020-11-23 15:46:42 +00:00
using Elsa.Activities.AzureServiceBus.Services ;
2020-11-22 21:16:28 +00:00
using Elsa.ActivityResults ;
using Elsa.Attributes ;
using Elsa.Serialization ;
using Elsa.Services ;
using Elsa.Services.Models ;
using Microsoft.Azure.ServiceBus ;
2020-11-23 15:46:42 +00:00
namespace Elsa.Activities.AzureServiceBus
2020-11-22 21:16:28 +00:00
{
[Trigger(Category = "Azure Service Bus", DisplayName = "Send Service Bus Message", Description = "Sends a message to the specified queue", Outcomes = new[] { OutcomeNames . Done } ) ]
public class SendAzureServiceBusMessage : Activity
{
2020-11-23 15:46:42 +00:00
private readonly IMessageSenderFactory _messageSenderFactory ;
2020-11-22 21:16:28 +00:00
private readonly IContentSerializer _serializer ;
2020-11-23 15:46:42 +00:00
public SendAzureServiceBusMessage ( IMessageSenderFactory messageSenderFactory , IContentSerializer serializer )
2020-11-22 21:16:28 +00:00
{
2020-11-23 15:46:42 +00:00
_messageSenderFactory = messageSenderFactory ;
2020-11-22 21:16:28 +00:00
_serializer = serializer ;
}
[ActivityProperty] public string QueueName { get ; set ; } = default ! ;
[ActivityProperty] public object Message { get ; set ; } = default ! ;
2020-12-02 21:28:36 +00:00
protected override async ValueTask < IActivityExecutionResult > OnExecuteAsync ( ActivityExecutionContext context )
2020-11-22 21:16:28 +00:00
{
2020-12-02 21:28:36 +00:00
var sender = await _messageSenderFactory . GetSenderAsync ( QueueName , context . CancellationToken ) ;
2020-11-22 21:16:28 +00:00
var json = _serializer . Serialize ( Message ) ;
var bytes = Encoding . UTF8 . GetBytes ( json ) ;
var message = new Message ( bytes ) ;
if ( ! string . IsNullOrWhiteSpace ( context . WorkflowExecutionContext . CorrelationId ) )
message . CorrelationId = context . WorkflowExecutionContext . CorrelationId ;
await sender . SendAsync ( message ) ;
return Done ( ) ;
}
}
}