Refactor code for readability and improved logging

Code in several activity files and Worker.cs that involve creating bookmarks and handling messages received was refactored to improve readability and clarity. The instantiated "CreateBookmarkArgs" parameter was split into its individual arguments for improved readability. Logging in the Worker.cs file was also fine-tuned, with explicit count of workflows triggered by service bus messages being included in the debug log.
This commit is contained in:
Sipke Schoorstra 2023-12-18 20:20:20 +01:00
parent 32e9cbfaa5
commit 5e95daaa7c
5 changed files with 40 additions and 27 deletions

View file

@ -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;
}
}

View file

@ -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;
/// </summary>
public class Worker : IAsyncDisposable
{
private static readonly string BookmarkName = TypeNameHelper.GenerateTypeName<MessageReceived>();
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
/// <summary>
/// Initializes a new instance of the <see cref="Worker"/> class.
/// </summary>
public Worker(string queueOrTopic, string? subscription, IWorkflowDispatcher workflowDispatcher, ServiceBusClient client, IHasher hasher, IWorkflowInbox workflowInbox, ILogger<Worker> logger)
public Worker(string queueOrTopic, string? subscription, ServiceBusClient client, IWorkflowInbox workflowInbox, ILogger<Worker> 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.
/// </summary>
public string QueueOrTopic { get; }
/// <summary>
/// The name of the subscription that this worker is processing. Only valid if the worker is processing a topic.
/// </summary>
@ -73,22 +68,22 @@ public class Worker : IAsyncDisposable
/// </summary>
/// <param name="cancellationToken">The cancellation token.</param>
public async Task StartAsync(CancellationToken cancellationToken = default) => await _processor.StartProcessingAsync(cancellationToken);
/// <summary>
/// Increments the ref count.
/// </summary>
public void IncrementRefCount() => RefCount++;
/// <summary>
/// Decrements the ref count.
/// </summary>
public void DecrementRefCount() => RefCount--;
/// <summary>
/// Disposes the worker.
/// </summary>
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<string, object> { [MessageReceived.InputKey] = messageModel };
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<MessageReceived>();
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);
}
}

View file

@ -30,11 +30,17 @@ public class CallAnswered : Activity<CallAnsweredPayload>
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
});
}
}

View file

@ -30,11 +30,17 @@ public class CallHangup : Activity<CallHangupPayload>
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
});
}
}

View file

@ -20,18 +20,18 @@ namespace Elsa.Telnyx.Activities;
public class WebhookEvent : Activity<Payload>
{
/// <inheritdoc />
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)
{
}
/// <inheritdoc />
public WebhookEvent(string eventType, string activityTypeName, Variable<Payload> result, int version = 1, [CallerFilePath]string? source = default, [CallerLineNumber]int? line = default)
public WebhookEvent(string eventType, string activityTypeName, Variable<Payload> result, int version = 1, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
: base(activityTypeName, version, source, line)
{
EventType = eventType;
Result = new(result);
}
/// <summary>
/// The Telnyx webhook event type to listen for.
/// </summary>
@ -47,8 +47,13 @@ public class WebhookEvent : Activity<Payload>
{
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
});
}
}