Configure Azure Service Bus activities
This commit is contained in:
parent
2e9617fdce
commit
570a07e656
|
|
@ -3,6 +3,7 @@ using Elsa.Activities.AzureServiceBus.Extensions;
|
|||
using Elsa.Activities.AzureServiceBus.Models;
|
||||
using Elsa.ActivityResults;
|
||||
using Elsa.Attributes;
|
||||
using Elsa.Expressions;
|
||||
using Elsa.Serialization;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
|
|
@ -19,8 +20,11 @@ namespace Elsa.Activities.AzureServiceBus
|
|||
_serializer = serializer;
|
||||
}
|
||||
|
||||
[ActivityProperty] public string QueueName { get; set; } = default!;
|
||||
[ActivityProperty] public Type MessageType { get; set; } = default!;
|
||||
[ActivityProperty(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
|
||||
public string QueueName { get; set; } = default!;
|
||||
|
||||
[ActivityProperty]
|
||||
public Type MessageType { get; set; } = default!;
|
||||
|
||||
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => context.WorkflowExecutionContext.IsFirstPass ? ExecuteInternal(context) : Suspend();
|
||||
protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) => ExecuteInternal(context);
|
||||
|
|
|
|||
|
|
@ -3,10 +3,13 @@ using Elsa.Activities.AzureServiceBus.Extensions;
|
|||
using Elsa.Activities.AzureServiceBus.Models;
|
||||
using Elsa.ActivityResults;
|
||||
using Elsa.Attributes;
|
||||
using Elsa.Expressions;
|
||||
using Elsa.Serialization;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
|
||||
// ReSharper disable ExplicitCallerInfoArgument
|
||||
// ReSharper disable once CheckNamespace
|
||||
namespace Elsa.Activities.AzureServiceBus
|
||||
{
|
||||
[Trigger(Category = "Azure Service Bus", DisplayName = "Service Bus Topic Message Received", Description = "Triggered when a message is received on the specified topic/subscription", Outcomes = new[] { OutcomeNames.Done })]
|
||||
|
|
@ -19,16 +22,21 @@ namespace Elsa.Activities.AzureServiceBus
|
|||
_serializer = serializer;
|
||||
}
|
||||
|
||||
[ActivityProperty] public string TopicName { get; set; } = default!;
|
||||
[ActivityProperty] public string SubscriptionName { get; set; } = default!;
|
||||
[ActivityProperty] public Type MessageType { get; set; } = default!;
|
||||
[ActivityProperty(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
|
||||
public string TopicName { get; set; } = default!;
|
||||
|
||||
[ActivityProperty(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
|
||||
public string SubscriptionName { get; set; } = default!;
|
||||
|
||||
[ActivityProperty]
|
||||
public Type MessageType { get; set; } = default!;
|
||||
|
||||
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => context.WorkflowExecutionContext.IsFirstPass ? ExecuteInternal(context) : Suspend();
|
||||
protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) => ExecuteInternal(context);
|
||||
|
||||
private IActivityExecutionResult ExecuteInternal(ActivityExecutionContext context)
|
||||
{
|
||||
var message = (MessageModel)context.Input!;
|
||||
var message = (MessageModel) context.Input!;
|
||||
var model = message.ReadBody(MessageType, _serializer);
|
||||
|
||||
return Done(model);
|
||||
|
|
@ -3,6 +3,7 @@ using System.Runtime.CompilerServices;
|
|||
using Elsa.Builders;
|
||||
|
||||
// ReSharper disable ExplicitCallerInfoArgument
|
||||
// ReSharper disable once CheckNamespace
|
||||
namespace Elsa.Activities.AzureServiceBus
|
||||
{
|
||||
public static class AzureServiceBusTopicMessageReceivedBuilderExtensions
|
||||
|
|
@ -2,6 +2,7 @@
|
|||
using Elsa.Activities.AzureServiceBus.Services;
|
||||
using Elsa.ActivityResults;
|
||||
using Elsa.Attributes;
|
||||
using Elsa.Expressions;
|
||||
using Elsa.Serialization;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
|
|
@ -20,13 +21,15 @@ namespace Elsa.Activities.AzureServiceBus
|
|||
_serializer = serializer;
|
||||
}
|
||||
|
||||
[ActivityProperty] public string QueueName { get; set; } = default!;
|
||||
[ActivityProperty] public object Message { get; set; } = default!;
|
||||
[ActivityProperty(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
|
||||
public string QueueName { get; set; } = default!;
|
||||
|
||||
[ActivityProperty(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 _messageSenderFactory.GetSenderAsync(QueueName, context.CancellationToken);
|
||||
|
||||
var message = Extensions.MessageBodyExtensions.CreateMessage(_serializer, Message);
|
||||
|
||||
if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId))
|
||||
|
|
|
|||
|
|
@ -13,22 +13,22 @@ namespace Elsa.Activities.AzureServiceBus.Extensions
|
|||
public static object ReadBody(this MessageModel message, Type type, IContentSerializer serializer)
|
||||
{
|
||||
if (type == typeof(string))
|
||||
return UTF8Encoding.UTF8.GetString(message.Body);
|
||||
return Encoding.UTF8.GetString(message.Body);
|
||||
|
||||
var bytes = message.Body;
|
||||
var json = Encoding.UTF8.GetString(bytes);
|
||||
return serializer.Deserialize(json, type)!;
|
||||
}
|
||||
|
||||
public static Message CreateMessage(IContentSerializer serializer, object Message)
|
||||
public static Message CreateMessage(IContentSerializer serializer, object message)
|
||||
{
|
||||
byte[] messageBytes;
|
||||
|
||||
if (Message.GetType() == typeof(string))
|
||||
messageBytes = UTF8Encoding.UTF8.GetBytes(Message as string);
|
||||
if (message is string s)
|
||||
messageBytes = Encoding.UTF8.GetBytes(s);
|
||||
else
|
||||
{
|
||||
var json = serializer.Serialize(Message);
|
||||
var json = serializer.Serialize(message);
|
||||
messageBytes = Encoding.UTF8.GetBytes(json);
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue