From ed47dcd600db0ebf307c7d00a0f608fa671d6171 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 8 Jun 2021 10:34:54 +0200 Subject: [PATCH] Update Azure Service Bus Send activities to use NodaTime primitives --- .../AzureServiceBusSendActivity.cs | 74 +++++++++++-------- .../SendAzureServiceBusQueueMessage.cs | 16 +--- .../SendAzureServiceBusTopicMessage.cs | 21 ++---- 3 files changed, 56 insertions(+), 55 deletions(-) diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusSendActivity.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusSendActivity.cs index 4cbe00957..330ab070e 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusSendActivity.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/AzureServiceBusSendActivity.cs @@ -11,16 +11,15 @@ using System; using System.Collections.Generic; using System.Text; using System.Threading.Tasks; +using NodaTime; namespace Elsa.Activities.AzureServiceBus { public abstract class AzureServiceBusSendActivity : Activity { - protected ISenderClient Sender; + private readonly IContentSerializer _serializer; - protected IContentSerializer _serializer; - - public AzureServiceBusSendActivity(IContentSerializer serializer) + protected AzureServiceBusSendActivity(IContentSerializer serializer) { _serializer = serializer; } @@ -28,34 +27,49 @@ namespace Elsa.Activities.AzureServiceBus [ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid, SyntaxNames.Json })] public object Message { get; set; } = default!; - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? CorrelationId { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string ContentType { get; set; } = default!; - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? Label { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? To { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] - public string MessageId { get; set; } = default!; - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + public string? MessageId { get; set; } + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? PartitionKey { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? ViaPartitionKey { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? ReplyTo { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? SessionId { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] - public DateTime ExpiresAtUtc { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] - public TimeSpan TimeToLive { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + public Instant ExpiresAtUtc { get; set; } + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + public Duration? TimeToLive { get; set; } + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string? ReplyToSessionId { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] - public DateTime ScheduledEnqueueTimeUtc { get; set; } - [ActivityInput(Category = PropertyCategories.Advanced, DefaultSyntax = SyntaxNames.Json, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid,SyntaxNames.Json },UIHint = ActivityInputUIHints.MultiLine)] - public IDictionary UserProperties { get; set; } = new Dictionary(); + + [ActivityInput(Category = PropertyCategories.Advanced, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + public Instant? ScheduledEnqueueTimeUtc { get; set; } + + [ActivityInput(Category = PropertyCategories.Advanced, DefaultSyntax = SyntaxNames.Json, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid, SyntaxNames.Json }, UIHint = ActivityInputUIHints.MultiLine)] + public IDictionary? UserProperties { get; set; } = new Dictionary(); + + protected abstract Task GetSenderAsync(); protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) { @@ -64,7 +78,8 @@ namespace Elsa.Activities.AzureServiceBus if (!string.IsNullOrWhiteSpace(context.WorkflowExecutionContext.CorrelationId)) message.CorrelationId = context.WorkflowExecutionContext.CorrelationId; - await Sender.SendAsync(message); + var sender = await GetSenderAsync(); + await sender.SendAsync(message); return Done(); } @@ -84,10 +99,12 @@ namespace Elsa.Activities.AzureServiceBus if (MessageId != null) message.MessageId = MessageId; - if (TimeToLive != null && TimeToLive > TimeSpan.Zero) - message.TimeToLive = TimeToLive; + + if (TimeToLive != null && TimeToLive > Duration.Zero) + message.TimeToLive = TimeToLive.Value.ToTimeSpan(); + if (ScheduledEnqueueTimeUtc != null) - message.ScheduledEnqueueTimeUtc = ScheduledEnqueueTimeUtc; + message.ScheduledEnqueueTimeUtc = ScheduledEnqueueTimeUtc.Value.ToDateTimeUtc(); if (UserProperties != null) foreach (var props in UserProperties) @@ -95,6 +112,5 @@ namespace Elsa.Activities.AzureServiceBus return message; } - } -} +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs index 2c8265f0d..a4f444bdf 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusQueueMessage.cs @@ -1,11 +1,9 @@ -using System.Threading.Tasks; +using System.Threading.Tasks; using Elsa.Activities.AzureServiceBus.Services; -using Elsa.ActivityResults; using Elsa.Attributes; using Elsa.Expressions; using Elsa.Serialization; -using Elsa.Services; -using Elsa.Services.Models; +using Microsoft.Azure.ServiceBus.Core; namespace Elsa.Activities.AzureServiceBus { @@ -15,18 +13,12 @@ namespace Elsa.Activities.AzureServiceBus private readonly IQueueMessageSenderFactory _queueMessageSenderFactory; public SendAzureServiceBusQueueMessage(IQueueMessageSenderFactory queueMessageSenderFactory, IContentSerializer serializer) - :base(serializer) - { + : base(serializer) => _queueMessageSenderFactory = queueMessageSenderFactory; - } [ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string QueueName { get; set; } = default!; - protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) - { - Sender = await _queueMessageSenderFactory.GetSenderAsync(QueueName, context.CancellationToken); - return await base.OnExecuteAsync(context); - } + protected override Task GetSenderAsync() => _queueMessageSenderFactory.GetSenderAsync(QueueName); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs index f14da10dc..f4e566f79 100644 --- a/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs +++ b/src/activities/Elsa.Activities.AzureServiceBus/Activities/SendAzureServiceBusMessage/SendAzureServiceBusTopicMessage.cs @@ -1,10 +1,9 @@ using System.Threading.Tasks; using Elsa.Activities.AzureServiceBus.Services; -using Elsa.ActivityResults; using Elsa.Attributes; +using Elsa.Expressions; using Elsa.Serialization; -using Elsa.Services; -using Elsa.Services.Models; +using Microsoft.Azure.ServiceBus.Core; namespace Elsa.Activities.AzureServiceBus { @@ -14,18 +13,12 @@ namespace Elsa.Activities.AzureServiceBus private readonly ITopicMessageSenderFactory _messageSenderFactory; public SendAzureServiceBusTopicMessage(ITopicMessageSenderFactory messageSenderFactory, IContentSerializer serializer) - :base(serializer) - { - _messageSenderFactory = messageSenderFactory; - } + : base(serializer) => + _messageSenderFactory = messageSenderFactory; - [ActivityInput] public string TopicName { get; set; } = default!; + [ActivityInput(SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + public string TopicName { get; set; } = default!; - protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) - { - Sender = await _messageSenderFactory.GetTopicSenderAsync(TopicName, context.CancellationToken); - - return await base.OnExecuteAsync(context); - } + protected override Task GetSenderAsync() => _messageSenderFactory.GetTopicSenderAsync(TopicName); } } \ No newline at end of file