diff --git a/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs b/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs index 69d6454c7..0ab8fe771 100644 --- a/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs +++ b/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs @@ -96,7 +96,6 @@ public class MessageReceived : Trigger { // Create bookmarks for when we receive the expected HTTP request. context.CreateBookmark(GetBookmarkPayload(context.ExpressionExecutionContext), Resume,false); - return; } } diff --git a/src/modules/Elsa.AzureServiceBus/Services/Worker.cs b/src/modules/Elsa.AzureServiceBus/Services/Worker.cs index 0e80f3164..49c91b531 100644 --- a/src/modules/Elsa.AzureServiceBus/Services/Worker.cs +++ b/src/modules/Elsa.AzureServiceBus/Services/Worker.cs @@ -4,7 +4,7 @@ using Elsa.AzureServiceBus.Models; using Elsa.Workflows.Core.Contracts; using Elsa.Workflows.Core.Helpers; using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Models; using Microsoft.Extensions.Logging; namespace Elsa.AzureServiceBus.Services; @@ -15,10 +15,7 @@ namespace Elsa.AzureServiceBus.Services; /// public class Worker : IAsyncDisposable { - private static readonly string BookmarkName = TypeNameHelper.GenerateTypeName(); private readonly ServiceBusProcessor _processor; - private readonly IWorkflowDispatcher _workflowDispatcher; - private readonly IHasher _hasher; private readonly IWorkflowInbox _workflowInbox; private readonly ILogger _logger; private int _refCount = 1; @@ -26,12 +23,10 @@ public class Worker : IAsyncDisposable /// /// Initializes a new instance of the class. /// - public Worker(string queueOrTopic, string? subscription, IWorkflowDispatcher workflowDispatcher, ServiceBusClient client, IHasher hasher, IWorkflowInbox workflowInbox, ILogger logger) + public Worker(string queueOrTopic, string? subscription, ServiceBusClient client, IWorkflowInbox workflowInbox, ILogger logger) { QueueOrTopic = queueOrTopic; Subscription = subscription == "" ? default : subscription; - _workflowDispatcher = workflowDispatcher; - _hasher = hasher; _workflowInbox = workflowInbox; _logger = logger; @@ -47,7 +42,7 @@ public class Worker : IAsyncDisposable /// The name of the queue or topic that this worker is processing. /// public string QueueOrTopic { get; } - + /// /// The name of the subscription that this worker is processing. Only valid if the worker is processing a topic. /// @@ -73,22 +68,22 @@ public class Worker : IAsyncDisposable /// /// The cancellation token. public async Task StartAsync(CancellationToken cancellationToken = default) => await _processor.StartProcessingAsync(cancellationToken); - + /// /// Increments the ref count. /// public void IncrementRefCount() => RefCount++; - + /// /// Decrements the ref count. /// public void DecrementRefCount() => RefCount--; - + /// /// Disposes the worker. /// public async ValueTask DisposeAsync() => await _processor.DisposeAsync(); - + private async Task OnMessageReceivedAsync(ProcessMessageEventArgs args) => await InvokeWorkflowsAsync(args.Message, args.CancellationToken); private Task OnErrorAsync(ProcessErrorEventArgs args) @@ -105,19 +100,20 @@ public class Worker : IAsyncDisposable var input = new Dictionary { [MessageReceived.InputKey] = messageModel }; var activityTypeName = ActivityTypeNameHelper.GenerateTypeName(); - var results = await _workflowInbox.SubmitAsync(new Workflows.Runtime.Models.NewWorkflowInboxMessage() + var results = await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage { ActivityTypeName = activityTypeName, BookmarkPayload = payload, Input = input, CorrelationId = correlationId - }); + }, cancellationToken); - _logger.LogInformation($"{results.WorkflowExecutionResults.Count()} workflow triggered by the service bus message"); + _logger.LogDebug("{Count} workflow triggered by the service bus message", results.WorkflowExecutionResults.Count); } - private ReceivedServiceBusMessageModel CreateMessageModel(ServiceBusReceivedMessage message) => - new( + private static ReceivedServiceBusMessageModel CreateMessageModel(ServiceBusReceivedMessage message) + { + return new( message.Body.ToArray(), message.Subject, message.ContentType, @@ -142,4 +138,5 @@ public class Worker : IAsyncDisposable message.DeadLetterSource, message.DeadLetterErrorDescription, message.ApplicationProperties); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs b/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs index 6940d30b2..54b9d1272 100644 --- a/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs +++ b/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs @@ -30,11 +30,17 @@ public class CallAnswered : Activity protected override void Execute(ActivityExecutionContext context) { var callControlIds = CallControlIds.Get(context); - + foreach (var callControlId in callControlIds) { var payload = new CallAnsweredBookmarkPayload(callControlId); - context.CreateBookmark(new CreateBookmarkArgs(payload, Resume, Type, IncludeActivityInstanceId: false)); + context.CreateBookmark(new CreateBookmarkArgs + { + Payload = payload, + Callback = Resume, + BookmarkName = Type, + IncludeActivityInstanceId = false + }); } } diff --git a/src/modules/Elsa.Telnyx/Activities/CallHangup.cs b/src/modules/Elsa.Telnyx/Activities/CallHangup.cs index 5131dc560..dc71b7219 100644 --- a/src/modules/Elsa.Telnyx/Activities/CallHangup.cs +++ b/src/modules/Elsa.Telnyx/Activities/CallHangup.cs @@ -30,11 +30,17 @@ public class CallHangup : Activity protected override void Execute(ActivityExecutionContext context) { var callControlIds = CallControlIds.Get(context); - + foreach (var callControlId in callControlIds) { var payload = new CallHangupBookmarkPayload(callControlId); - context.CreateBookmark(new CreateBookmarkArgs(payload, Resume, Type, IncludeActivityInstanceId: false)); + context.CreateBookmark(new CreateBookmarkArgs + { + Payload = payload, + Callback = Resume, + BookmarkName = Type, + IncludeActivityInstanceId = false + }); } } diff --git a/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs b/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs index 03fadd483..819cc20d4 100644 --- a/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs +++ b/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs @@ -20,18 +20,18 @@ namespace Elsa.Telnyx.Activities; public class WebhookEvent : Activity { /// - public WebhookEvent([CallerFilePath]string? source = default, [CallerLineNumber]int? line = default) : base(source, line) + public WebhookEvent([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) { } /// - public WebhookEvent(string eventType, string activityTypeName, Variable result, int version = 1, [CallerFilePath]string? source = default, [CallerLineNumber]int? line = default) + public WebhookEvent(string eventType, string activityTypeName, Variable result, int version = 1, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(activityTypeName, version, source, line) { EventType = eventType; Result = new(result); } - + /// /// The Telnyx webhook event type to listen for. /// @@ -47,8 +47,13 @@ public class WebhookEvent : Activity { var eventType = EventType; var payload = new WebhookEventBookmarkPayload(eventType); - - context.CreateBookmark(new CreateBookmarkArgs(payload, Resume, Type, false)); + + context.CreateBookmark(new CreateBookmarkArgs + { + Payload = payload, + Callback = Resume, + BookmarkName = Type + }); } }