From f6cb6cc055a7190fa0f2d974275453a236fbb162 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 22 Jan 2021 14:58:58 +0100 Subject: [PATCH] Reorder startup tasks --- .../Consumers/RunWorkflowInstanceConsumer.cs | 1 - .../ElsaServiceCollectionExtensions.cs | 3 ++- .../CreateSubscriptions.cs | 11 ++++++---- .../StartupRunnerHostedService.cs | 2 +- .../Runtime/ServiceCollectionExtensions.cs | 1 + .../Elsa.Core/Services/BackgroundWorker.cs | 6 ++--- .../Elsa.Core/Services/ServiceBusFactory.cs | 5 +++-- .../Extensions/ServiceCollectionExtensions.cs | 3 +-- .../StartupTasks/RunEFCoreMigrations.cs | 22 ++++++++++++++----- 9 files changed, 34 insertions(+), 20 deletions(-) rename src/core/Elsa.Core/{StartupTasks => HostedServices}/CreateSubscriptions.cs (72%) rename src/core/Elsa.Core/{Runtime => HostedServices}/StartupRunnerHostedService.cs (96%) diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs index 7ffa22dd0..f46179b71 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs @@ -47,7 +47,6 @@ namespace Elsa.Consumers { // Reschedule message. _logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Rescheduling message", workflowInstanceId); - await Task.Delay(TimeSpan.FromSeconds(1)); await _commandSender.SendAsync(message); return; } diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index d22c57912..0cdd22b9b 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -53,7 +53,6 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton(options.DistributedLockProviderFactory) .AddSingleton(options.SignalFactory) .AddSingleton(options.StorageFactory) - .AddStartupTask() .AddStartupTask(); options @@ -68,6 +67,8 @@ namespace Microsoft.Extensions.DependencyInjection services.Decorate(); services.Decorate(); services.Decorate(); + + services.AddHostedService(); return services; } diff --git a/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs b/src/core/Elsa.Core/HostedServices/CreateSubscriptions.cs similarity index 72% rename from src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs rename to src/core/Elsa.Core/HostedServices/CreateSubscriptions.cs index 2fb0fd7ca..467178a91 100644 --- a/src/core/Elsa.Core/StartupTasks/CreateSubscriptions.cs +++ b/src/core/Elsa.Core/HostedServices/CreateSubscriptions.cs @@ -3,10 +3,11 @@ using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.Services; +using Microsoft.Extensions.Hosting; -namespace Elsa.StartupTasks +namespace Elsa.HostedServices { - public class CreateSubscriptions : IStartupTask + public class CreateSubscriptions : IHostedService { private readonly IServiceBusFactory _serviceBusFactory; private readonly IEnumerable _messageTypes; @@ -16,8 +17,8 @@ namespace Elsa.StartupTasks _serviceBusFactory = serviceBusFactory; _messageTypes = elsaOptions.MessageTypes; } - - public async Task ExecuteAsync(CancellationToken cancellationToken = default) + + public async Task StartAsync(CancellationToken cancellationToken) { foreach (var messageType in _messageTypes) { @@ -25,5 +26,7 @@ namespace Elsa.StartupTasks await bus.Subscribe(messageType); } } + + public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Runtime/StartupRunnerHostedService.cs b/src/core/Elsa.Core/HostedServices/StartupRunnerHostedService.cs similarity index 96% rename from src/core/Elsa.Core/Runtime/StartupRunnerHostedService.cs rename to src/core/Elsa.Core/HostedServices/StartupRunnerHostedService.cs index cd8ba8c5b..6c1d3809d 100644 --- a/src/core/Elsa.Core/Runtime/StartupRunnerHostedService.cs +++ b/src/core/Elsa.Core/HostedServices/StartupRunnerHostedService.cs @@ -5,7 +5,7 @@ using Elsa.Services; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; -namespace Elsa.Runtime +namespace Elsa.HostedServices { public class StartupRunnerHostedService : IHostedService { diff --git a/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs b/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs index f9ce7e5b4..ab38b911c 100644 --- a/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs @@ -1,4 +1,5 @@ using System; +using Elsa.HostedServices; using Elsa.Services; using Microsoft.Extensions.DependencyInjection; diff --git a/src/core/Elsa.Core/Services/BackgroundWorker.cs b/src/core/Elsa.Core/Services/BackgroundWorker.cs index 19111071a..93ffdd414 100644 --- a/src/core/Elsa.Core/Services/BackgroundWorker.cs +++ b/src/core/Elsa.Core/Services/BackgroundWorker.cs @@ -14,7 +14,7 @@ namespace Elsa.Services public async Task ScheduleTask(string channelName, Func task, CancellationToken cancellationToken, Duration? channelTtl) { - var channel = GetOrCreateChannel(channelName, channelTtl, cancellationToken); + var channel = await GetOrCreateChannel(channelName, channelTtl, cancellationToken); await channel.Writer.WriteAsync(task, cancellationToken); } @@ -27,9 +27,9 @@ namespace Elsa.Services }, cancellationToken, channelTtl); } - private Channel> GetOrCreateChannel(string channelName, Duration? channelTtl, CancellationToken cancellationToken) + private async Task>> GetOrCreateChannel(string channelName, Duration? channelTtl, CancellationToken cancellationToken) { - _semaphore.Wait(cancellationToken); + await _semaphore.WaitAsync(cancellationToken); try { diff --git a/src/core/Elsa.Core/Services/ServiceBusFactory.cs b/src/core/Elsa.Core/Services/ServiceBusFactory.cs index c47efd955..774b89cc1 100644 --- a/src/core/Elsa.Core/Services/ServiceBusFactory.cs +++ b/src/core/Elsa.Core/Services/ServiceBusFactory.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Concurrent; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; @@ -12,7 +13,7 @@ namespace Elsa.Services { private readonly ElsaOptions _elsaOptions; private readonly IServiceProvider _serviceProvider; - private readonly IDictionary _serviceBuses = new Dictionary(); + private readonly ConcurrentDictionary _serviceBuses = new(); private readonly DependencyInjectionHandlerActivator _handlerActivator; private readonly SemaphoreSlim _semaphore = new(1); @@ -40,7 +41,7 @@ namespace Elsa.Services _elsaOptions.ConfigureServiceBusEndpoint(configureContext); var newBus = configurer.Start(); - _serviceBuses.Add(queueName, newBus); + _serviceBuses.TryAdd(queueName, newBus); return newBus; } diff --git a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Extensions/ServiceCollectionExtensions.cs b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Extensions/ServiceCollectionExtensions.cs index 37910b9e7..815eb59ff 100644 --- a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Extensions/ServiceCollectionExtensions.cs +++ b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Extensions/ServiceCollectionExtensions.cs @@ -1,7 +1,6 @@ using System; using Elsa.Persistence.EntityFramework.Core.StartupTasks; using Elsa.Persistence.EntityFramework.Core.Stores; -using Elsa.Runtime; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; @@ -21,7 +20,7 @@ namespace Elsa.Persistence.EntityFramework.Core.Extensions .AddScoped(); if (autorunMigrations) - elsa.Services.AddStartupTask(); + elsa.Services.AddHostedService(); return elsa .UseWorkflowDefinitionStore(sp => sp.GetRequiredService()) diff --git a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/StartupTasks/RunEFCoreMigrations.cs b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/StartupTasks/RunEFCoreMigrations.cs index 0ba030a8f..eb8e70ae8 100644 --- a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/StartupTasks/RunEFCoreMigrations.cs +++ b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/StartupTasks/RunEFCoreMigrations.cs @@ -1,17 +1,27 @@ -using System.Threading; +using System; +using System.Threading; using System.Threading.Tasks; -using Elsa.Services; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; namespace Elsa.Persistence.EntityFramework.Core.StartupTasks { /// /// Executes EF Core migrations. /// - public class RunEFCoreMigrations : IStartupTask + public class RunEFCoreMigrations : IHostedService { - private readonly ElsaContext _elsaContext; - public RunEFCoreMigrations(ElsaContext elsaContext) => _elsaContext = elsaContext; - public async Task ExecuteAsync(CancellationToken cancellationToken = default) => await _elsaContext.Database.MigrateAsync(cancellationToken); + private readonly IServiceProvider _serviceProvider; + public RunEFCoreMigrations(IServiceProvider serviceProvider) => _serviceProvider = serviceProvider; + + public async Task StartAsync(CancellationToken cancellationToken) + { + using var scope = _serviceProvider.CreateScope(); + var elsaContext = scope.ServiceProvider.GetRequiredService(); + await elsaContext.Database.MigrateAsync(cancellationToken); + } + + public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; } } \ No newline at end of file