From 49d4e5bb3656336a15fbb9580081276a6dc7d590 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 10 Jan 2021 13:43:34 +0100 Subject: [PATCH] Optimize workflow trigger re-indexing --- .../Jobs/RunQuartzWorkflowJob.cs | 7 ++- .../Triggers/IWorkflowSelector.cs | 2 +- .../Consumers/RunWorkflowInstanceConsumer.cs | 10 +++- src/core/Elsa.Core/ElsaOptions.cs | 9 +--- .../RebusOptionsConfigurerExtensions.cs | 16 ++++++ .../Handlers/LogWorkflowExecution.cs | 46 ++++++++++++++++ .../Elsa.Core/Handlers/PersistWorkflow.cs | 7 +-- src/core/Elsa.Core/Handlers/UpdateTriggers.cs | 8 +-- src/core/Elsa.Core/ServiceBusOptions.cs | 13 ++++- src/core/Elsa.Core/Services/WorkflowRunner.cs | 6 +-- .../ResumeRunningWorkflowsTask.cs | 24 ++++++--- .../Elsa.Core/Triggers/WorkflowSelector.cs | 52 +++++++++++++------ .../Extensions/ElsaOptionsExtensions.cs | 7 +-- 13 files changed, 154 insertions(+), 53 deletions(-) create mode 100644 src/core/Elsa.Core/Extensions/RebusOptionsConfigurerExtensions.cs create mode 100644 src/core/Elsa.Core/Handlers/LogWorkflowExecution.cs diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs b/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs index a2da6a7b5..e8b1e2198 100644 --- a/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs +++ b/src/activities/Elsa.Activities.Timers.Quartz/Jobs/RunQuartzWorkflowJob.cs @@ -1,5 +1,4 @@ -using System; -using System.Diagnostics; +using System.Diagnostics; using System.Threading.Tasks; using Elsa.DistributedLock; using Elsa.Models; @@ -18,12 +17,12 @@ namespace Elsa.Activities.Timers.Quartz.Jobs private readonly IWorkflowQueue _workflowQueue; private readonly IDistributedLockProvider _distributedLockProvider; private readonly ILogger _logger; - private Stopwatch _stopwatch = new(); + private readonly Stopwatch _stopwatch = new(); public RunQuartzWorkflowJob( IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, - IWorkflowQueue workflowQueue, + IWorkflowQueue workflowQueue, IDistributedLockProvider distributedLockProvider, ILogger logger) { diff --git a/src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs b/src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs index 99cb1dfbb..dc682119b 100644 --- a/src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs +++ b/src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs @@ -15,7 +15,7 @@ namespace Elsa.Triggers Func evaluate, CancellationToken cancellationToken = default); - Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default); + Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, string? workflowInstanceId, CancellationToken cancellationToken = default); Task RemoveTriggerAsync(ITrigger trigger, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs index a059b9f9c..b46d616b2 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs @@ -1,4 +1,5 @@ using System; +using System.Diagnostics; using System.Threading.Tasks; using Elsa.DistributedLock; using Elsa.Messages; @@ -17,6 +18,7 @@ namespace Elsa.Consumers private readonly IDistributedLockProvider _distributedLockProvider; private readonly IEventPublisher _eventPublisher; private readonly ILogger _logger; + private readonly Stopwatch _stopwatch = new(); public RunWorkflowInstanceConsumer( IWorkflowRunner workflowRunner, @@ -35,10 +37,12 @@ namespace Elsa.Consumers public async Task Handle(RunWorkflowInstance message) { var workflowInstanceId = message.WorkflowInstanceId; + var lockKey = workflowInstanceId; _logger.LogDebug("Acquiring lock on workflow instance {WorkflowInstanceId}.", workflowInstanceId); + _stopwatch.Restart(); - if (!await _distributedLockProvider.AcquireLockAsync(workflowInstanceId)) + if (!await _distributedLockProvider.AcquireLockAsync(lockKey)) { // Reschedule message. _logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Rescheduling message.", workflowInstanceId); @@ -61,7 +65,9 @@ namespace Elsa.Consumers } finally { - await _distributedLockProvider.ReleaseLockAsync(workflowInstanceId); + await _distributedLockProvider.ReleaseLockAsync(lockKey); + _stopwatch.Stop(); + _logger.LogDebug("Held lock on workflow instance {WorkflowInstanceId} for {ElapsedTime}.", workflowInstanceId, _stopwatch.Elapsed); } } diff --git a/src/core/Elsa.Core/ElsaOptions.cs b/src/core/Elsa.Core/ElsaOptions.cs index ff3b32c8e..c0ab60c1e 100644 --- a/src/core/Elsa.Core/ElsaOptions.cs +++ b/src/core/Elsa.Core/ElsaOptions.cs @@ -237,14 +237,7 @@ namespace Elsa .Subscriptions(s => s.StoreInMemory(store)) .Transport(t => t.UseInMemoryTransport(transport, queueName)) .Routing(r => r.TypeBased().Map(context.MessageTypeMap)) - .Options(options => - { - if(ServiceBusOptions.NumberOfWorkers != null) - options.SetNumberOfWorkers(ServiceBusOptions.NumberOfWorkers.Value); - - if(ServiceBusOptions.MaxParallelism != null) - options.SetMaxParallelism(ServiceBusOptions.MaxParallelism.Value); - }); + .Options(options => options.Apply(ServiceBusOptions)); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/RebusOptionsConfigurerExtensions.cs b/src/core/Elsa.Core/Extensions/RebusOptionsConfigurerExtensions.cs new file mode 100644 index 000000000..6ecdf4296 --- /dev/null +++ b/src/core/Elsa.Core/Extensions/RebusOptionsConfigurerExtensions.cs @@ -0,0 +1,16 @@ +using Rebus.Config; + +namespace Elsa +{ + public static class RebusOptionsConfigurerExtensions + { + public static void Apply(this OptionsConfigurer configurer, ServiceBusOptions options) + { + if(options.NumberOfWorkers != null) + configurer.SetNumberOfWorkers(options.NumberOfWorkers.Value); + + if(options.MaxParallelism != null) + configurer.SetMaxParallelism(options.MaxParallelism.Value); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Handlers/LogWorkflowExecution.cs b/src/core/Elsa.Core/Handlers/LogWorkflowExecution.cs new file mode 100644 index 000000000..065a78868 --- /dev/null +++ b/src/core/Elsa.Core/Handlers/LogWorkflowExecution.cs @@ -0,0 +1,46 @@ +using System.Diagnostics; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Events; +using MediatR; +using Microsoft.Extensions.Logging; + +namespace Elsa.Handlers +{ + public class LogWorkflowExecution : INotificationHandler, INotificationHandler, INotificationHandler + { + private readonly ILogger _logger; + private readonly Stopwatch _stopwatch = new(); + + public LogWorkflowExecution(ILogger logger) + { + _logger = logger; + } + + public Task Handle(ActivityExecuting notification, CancellationToken cancellationToken) + { + _stopwatch.Start(); + var activityBlueprint = notification.Activity; + var activityId = activityBlueprint.Id; + var workflowInstanceId = notification.WorkflowExecutionContext.WorkflowInstance.Id; + _logger.LogDebug("Executing activity {ActivityType} {ActivityId} for workflow {workflowInstanceId}", activityBlueprint.Type, activityId, workflowInstanceId); + return Task.CompletedTask; + } + + public Task Handle(ActivityExecuted notification, CancellationToken cancellationToken) + { + _stopwatch.Stop(); + var activityBlueprint = notification.Activity; + var activityId = activityBlueprint.Id; + var workflowInstanceId = notification.WorkflowExecutionContext.WorkflowInstance.Id; + _logger.LogDebug("Executed activity {ActivityType} {ActivityId} for workflow {workflowInstanceId} in {ElapsedTime}", activityBlueprint.Type, activityId, workflowInstanceId, _stopwatch.Elapsed); + return Task.CompletedTask; + } + + public Task Handle(WorkflowExecutionBurstCompleted notification, CancellationToken cancellationToken) + { + _logger.LogDebug("Burst of workflow execution completed for workflow {WorkflowInstanceId}", notification.WorkflowExecutionContext.WorkflowInstance.Id); + return Task.CompletedTask; + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Handlers/PersistWorkflow.cs b/src/core/Elsa.Core/Handlers/PersistWorkflow.cs index ab4aea77b..9baf66078 100644 --- a/src/core/Elsa.Core/Handlers/PersistWorkflow.cs +++ b/src/core/Elsa.Core/Handlers/PersistWorkflow.cs @@ -12,7 +12,7 @@ namespace Elsa.Handlers public class PersistWorkflow : INotificationHandler, INotificationHandler, - INotificationHandler, + INotificationHandler, INotificationHandler { private readonly IWorkflowInstanceStore _workflowInstanceStore; @@ -36,9 +36,9 @@ namespace Elsa.Handlers await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken); } - public async Task Handle(ActivityExecuted notification, CancellationToken cancellationToken) + public async Task Handle(WorkflowExecutionPassCompleted notification, CancellationToken cancellationToken) { - if (notification.WorkflowExecutionContext.WorkflowBlueprint.PersistenceBehavior == WorkflowPersistenceBehavior.ActivityExecuted || notification.Activity.PersistWorkflow) + if (notification.WorkflowExecutionContext.WorkflowBlueprint.PersistenceBehavior == WorkflowPersistenceBehavior.ActivityExecuted || notification.ActivityExecutionContext.ActivityBlueprint.PersistWorkflow) await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken); } @@ -61,6 +61,7 @@ namespace Elsa.Handlers { var workflowInstance = workflowExecutionContext.WorkflowInstance; await _workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken); + _logger.LogDebug("Committed workflow {WorkflowInstanceId} to storage", workflowInstance.Id); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Handlers/UpdateTriggers.cs b/src/core/Elsa.Core/Handlers/UpdateTriggers.cs index f7f8cf287..61fd01845 100644 --- a/src/core/Elsa.Core/Handlers/UpdateTriggers.cs +++ b/src/core/Elsa.Core/Handlers/UpdateTriggers.cs @@ -23,23 +23,23 @@ namespace Elsa.Handlers public async Task Handle(WorkflowInstanceSaved notification, CancellationToken cancellationToken) { var workflowInstance = notification.WorkflowInstance; - await UpdateTriggersAsync(workflowInstance.DefinitionId, workflowInstance.TenantId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken); + await UpdateTriggersAsync(workflowInstance.DefinitionId, workflowInstance.Id, workflowInstance.TenantId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken); } public async Task Handle(ManyWorkflowInstancesDeleted notification, CancellationToken cancellationToken) { var workflowInstance = notification.WorkflowInstances.First(); - await UpdateTriggersAsync(workflowInstance.DefinitionId, workflowInstance.TenantId, VersionOptions.Latest, cancellationToken); + await UpdateTriggersAsync(workflowInstance.DefinitionId, null, workflowInstance.TenantId, VersionOptions.Latest, cancellationToken); } - private async Task UpdateTriggersAsync(string workflowDefinitionId, string? tenantId, VersionOptions versionOptions, CancellationToken cancellationToken) + private async Task UpdateTriggersAsync(string workflowDefinitionId, string? workflowInstanceId, string? tenantId, VersionOptions versionOptions, CancellationToken cancellationToken) { var workflowBlueprint = await _workflowRegistry.GetWorkflowAsync(workflowDefinitionId, tenantId, versionOptions, cancellationToken); if (workflowBlueprint == null) return; - await _workflowSelector.UpdateTriggersAsync(workflowBlueprint, cancellationToken); + await _workflowSelector.UpdateTriggersAsync(workflowBlueprint, workflowInstanceId, cancellationToken); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/ServiceBusOptions.cs b/src/core/Elsa.Core/ServiceBusOptions.cs index 51f2d5812..ad029a348 100644 --- a/src/core/Elsa.Core/ServiceBusOptions.cs +++ b/src/core/Elsa.Core/ServiceBusOptions.cs @@ -1,8 +1,19 @@ -namespace Elsa +using Rebus.Config; + +namespace Elsa { public class ServiceBusOptions { public int? NumberOfWorkers { get; set; } public int? MaxParallelism { get; set; } + + public void Apply(OptionsConfigurer configurer) + { + if(NumberOfWorkers != null) + configurer.SetNumberOfWorkers(NumberOfWorkers.Value); + + if(MaxParallelism != null) + configurer.SetMaxParallelism(MaxParallelism.Value); + } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowRunner.cs b/src/core/Elsa.Core/Services/WorkflowRunner.cs index b32e14c83..81c649f37 100644 --- a/src/core/Elsa.Core/Services/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/WorkflowRunner.cs @@ -285,15 +285,13 @@ namespace Elsa.Services var activityBlueprint = workflowBlueprint.GetActivity(currentActivityId)!; var activityExecutionContext = new ActivityExecutionContext(scope, workflowExecutionContext, activityBlueprint, scheduledActivity.Input, cancellationToken); var activity = await activityExecutionContext.ActivateActivityAsync(cancellationToken); - _logger.LogDebug("Executing activity {ActivityType}: {ActivityId}", activityBlueprint.Type, currentActivityId); - var result = await activityOperation(activityExecutionContext, activity); await _mediator.Publish(new ActivityExecuting(activityExecutionContext), cancellationToken); + var result = await activityOperation(activityExecutionContext, activity); + await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken); await result.ExecuteAsync(activityExecutionContext, cancellationToken); workflowExecutionContext.WorkflowInstance.Output = activityExecutionContext.Output; workflowExecutionContext.ExecutionLog.Add(activity.Id); workflowExecutionContext.PruneActivityData(); - await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken); - activityOperation = Execute; workflowExecutionContext.CompletePass(); await _mediator.Publish(new WorkflowExecutionPassCompleted(workflowExecutionContext, activityExecutionContext), cancellationToken); diff --git a/src/core/Elsa.Core/StartupTasks/ResumeRunningWorkflowsTask.cs b/src/core/Elsa.Core/StartupTasks/ResumeRunningWorkflowsTask.cs index b1640a679..cd71159a9 100644 --- a/src/core/Elsa.Core/StartupTasks/ResumeRunningWorkflowsTask.cs +++ b/src/core/Elsa.Core/StartupTasks/ResumeRunningWorkflowsTask.cs @@ -1,3 +1,4 @@ +using System; using System.Threading; using System.Threading.Tasks; using Elsa.DistributedLock; @@ -30,15 +31,24 @@ namespace Elsa.StartupTasks public async Task ExecuteAsync(CancellationToken cancellationToken = default) { - if (!await _distributedLockProvider.AcquireLockAsync(GetType().Name, cancellationToken)) + var lockKey = GetType().Name; + + if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken)) return; - var instances = await _workflowInstanceStore.FindManyAsync(new WorkflowStatusSpecification(WorkflowStatus.Running), cancellationToken: cancellationToken); - - foreach (var instance in instances) - await _workflowScheduler.RunWorkflowAsync( - instance, - cancellationToken: cancellationToken); + try + { + var instances = await _workflowInstanceStore.FindManyAsync(new WorkflowStatusSpecification(WorkflowStatus.Running), cancellationToken: cancellationToken); + + foreach (var instance in instances) + await _workflowScheduler.RunWorkflowAsync( + instance, + cancellationToken: cancellationToken); + } + finally + { + await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken); + } } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Triggers/WorkflowSelector.cs b/src/core/Elsa.Core/Triggers/WorkflowSelector.cs index 0a16e34e1..cacb88355 100644 --- a/src/core/Elsa.Core/Triggers/WorkflowSelector.cs +++ b/src/core/Elsa.Core/Triggers/WorkflowSelector.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Diagnostics; using System.Linq; using System.Threading; using System.Threading.Tasks; @@ -10,6 +11,7 @@ using Elsa.Services; using Elsa.Services.Models; using Microsoft.Extensions.Caching.Memory; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; using Open.Linq.AsyncExtensions; namespace Elsa.Triggers @@ -24,7 +26,9 @@ namespace Elsa.Triggers private readonly IEnumerable _triggerProviders; private readonly IMemoryCache _memoryCache; private readonly IServiceProvider _serviceProvider; + private readonly ILogger _logger; private IDictionary>? _descriptors; + private readonly Stopwatch _stopwatch = new(); public WorkflowSelector( IWorkflowRegistry workflowRegistry, @@ -33,7 +37,8 @@ namespace Elsa.Triggers IWorkflowContextManager workflowContextManager, IEnumerable triggerProviders, IMemoryCache memoryCache, - IServiceProvider serviceProvider) + IServiceProvider serviceProvider, + ILogger logger) { _workflowRegistry = workflowRegistry; _workflowFactory = workflowFactory; @@ -42,6 +47,7 @@ namespace Elsa.Triggers _triggerProviders = triggerProviders; _memoryCache = memoryCache; _serviceProvider = serviceProvider; + _logger = logger; } public async Task> SelectWorkflowsAsync( @@ -62,11 +68,23 @@ namespace Elsa.Triggers return SelectWorkflowsAsync(_descriptors, triggerType, evaluate).Select(result => result.WorkflowBlueprint.Activities.First(x => x.Id == result.ActivityId)); } - public async Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default) + public async Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, string? workflowInstanceId, CancellationToken cancellationToken = default) { + if (workflowInstanceId == null) + _logger.LogDebug("Updating triggers for workflows of {WorkflowDefinitionId}", workflowBlueprint.Id); + else + _logger.LogDebug("Updating triggers for workflow {WorkflowInstanceId}", workflowInstanceId); + + _stopwatch.Restart(); _descriptors ??= await GetDescriptorsAsync(cancellationToken); - _descriptors[workflowBlueprint.Id] = await BuildDescriptorsForAsync(workflowBlueprint, cancellationToken).ToList(); + _descriptors[workflowBlueprint.Id] = await BuildDescriptorsForAsync(workflowBlueprint, workflowInstanceId, cancellationToken).ToList(); _memoryCache.Set(CacheKey, _descriptors); + _stopwatch.Stop(); + + if (workflowInstanceId == null) + _logger.LogDebug("Updated triggers for workflows of {WorkflowDefinitionId} in {ElapsedTime}", workflowBlueprint.Id, _stopwatch.Elapsed); + else + _logger.LogDebug("Updated triggers for workflow {WorkflowInstanceId} in {ElapsedTime}", workflowInstanceId, _stopwatch.Elapsed); } public async Task RemoveTriggerAsync(ITrigger trigger, CancellationToken cancellationToken = default) @@ -113,7 +131,7 @@ namespace Elsa.Triggers foreach (var workflowBlueprint in workflowBlueprints) { - var blueprintDescriptors = await BuildDescriptorsForAsync(workflowBlueprint, cancellationToken); + var blueprintDescriptors = await BuildDescriptorsForAsync(workflowBlueprint, null, cancellationToken); descriptors.AddRange(blueprintDescriptors); } @@ -122,7 +140,7 @@ namespace Elsa.Triggers .ToDictionary(x => x.Key, x => (ICollection) x.ToList()); } - private async Task> BuildDescriptorsForAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken) + private async Task> BuildDescriptorsForAsync(IWorkflowBlueprint workflowBlueprint, string? workflowInstanceId, CancellationToken cancellationToken) { var descriptors = new List(); @@ -131,10 +149,12 @@ namespace Elsa.Triggers descriptors.AddRange(startTriggers); // Build triggers for workflow instances. - var specification = new WorkflowInstanceDefinitionIdSpecification(workflowBlueprint.Id) - .WithTenant(workflowBlueprint.TenantId) - .WithStatus(WorkflowStatus.Suspended); - + var specification = workflowInstanceId == null + ? new WorkflowInstanceDefinitionIdSpecification(workflowBlueprint.Id) + .WithTenant(workflowBlueprint.TenantId) + .WithStatus(WorkflowStatus.Suspended) + : new WorkflowInstanceIdSpecification(workflowInstanceId); + var workflowInstances = await _workflowInstanceStore .FindManyAsync(specification, cancellationToken: cancellationToken) .ToList(); @@ -153,7 +173,7 @@ namespace Elsa.Triggers { var startActivities = workflowBlueprint.GetStartActivities(); var workflowInstance = await _workflowFactory.InstantiateAsync(workflowBlueprint, cancellationToken: cancellationToken); - + // This is a transient workflow instance; setting EntityId to null ensures trigger providers don't try and load the workflow instance. workflowInstance.Id = null!; return await BuildDescriptorsAsync(workflowBlueprint, startActivities, workflowInstance, cancellationToken); @@ -170,12 +190,12 @@ namespace Elsa.Triggers var scope = _serviceProvider.CreateScope(); var workflowExecutionContext = new WorkflowExecutionContext(scope, workflowBlueprint, workflowInstance); var isTransientWorkflowInstance = workflowInstance.Id == null!; - - workflowExecutionContext.WorkflowContext = - workflowBlueprint.ContextOptions != null && - !isTransientWorkflowInstance && - !string.IsNullOrWhiteSpace(workflowInstance.ContextId) - ? await _workflowContextManager.LoadContext(new LoadWorkflowContext(workflowExecutionContext), cancellationToken) + + workflowExecutionContext.WorkflowContext = + workflowBlueprint.ContextOptions != null && + !isTransientWorkflowInstance && + !string.IsNullOrWhiteSpace(workflowInstance.ContextId) + ? await _workflowContextManager.LoadContext(new LoadWorkflowContext(workflowExecutionContext), cancellationToken) : default; foreach (var blockingActivity in blockingActivities) diff --git a/src/servicebus/Elsa.Rebus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs b/src/servicebus/Elsa.Rebus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs index 7f92342b1..9c41227da 100644 --- a/src/servicebus/Elsa.Rebus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs +++ b/src/servicebus/Elsa.Rebus.AzureServiceBus/Extensions/ElsaOptionsExtensions.cs @@ -11,10 +11,10 @@ namespace Elsa.Rebus.AzureServiceBus { public static ElsaOptions UseAzureServiceBus(this ElsaOptions elsaOptions, string connectionString, ITokenProvider? tokenProvider = default) { - return elsaOptions.UseServiceBus(context => ConfigureAzureServiceBusEndpoint(context, connectionString, tokenProvider)); + return elsaOptions.UseServiceBus(context => ConfigureAzureServiceBusEndpoint(elsaOptions, context, connectionString, tokenProvider)); } - private static void ConfigureAzureServiceBusEndpoint(ServiceBusEndpointConfigurationContext context, string connectionString, ITokenProvider? tokenProvider) + private static void ConfigureAzureServiceBusEndpoint(ElsaOptions elsaOptions, ServiceBusEndpointConfigurationContext context, string connectionString, ITokenProvider? tokenProvider) { var queueName = context.QueueName; var loggerFactory = context.ServiceProvider.GetRequiredService(); @@ -22,7 +22,8 @@ namespace Elsa.Rebus.AzureServiceBus context.Configurer .Logging(l => l.MicrosoftExtensionsLogging(loggerFactory)) .Transport(t => t.UseAzureServiceBus(connectionString, queueName, tokenProvider)) - .Routing(r => r.TypeBased().Map(context.MessageTypeMap)); + .Routing(r => r.TypeBased().Map(context.MessageTypeMap)) + .Options(o => o.Apply(elsaOptions.ServiceBusOptions)); } } } \ No newline at end of file