Reorder startup tasks

This commit is contained in:
Sipke Schoorstra 2021-01-22 14:58:58 +01:00
parent 0a11cced56
commit f6cb6cc055
9 changed files with 34 additions and 20 deletions

View file

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

View file

@ -53,7 +53,6 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton(options.DistributedLockProviderFactory)
.AddSingleton(options.SignalFactory)
.AddSingleton(options.StorageFactory)
.AddStartupTask<CreateSubscriptions>()
.AddStartupTask<ContinueRunningWorkflowsTask>();
options
@ -68,6 +67,8 @@ namespace Microsoft.Extensions.DependencyInjection
services.Decorate<IWorkflowDefinitionStore, InitializingWorkflowDefinitionStore>();
services.Decorate<IWorkflowDefinitionStore, EventPublishingWorkflowDefinitionStore>();
services.Decorate<IWorkflowInstanceStore, EventPublishingWorkflowInstanceStore>();
services.AddHostedService<CreateSubscriptions>();
return services;
}

View file

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

View file

@ -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
{

View file

@ -1,4 +1,5 @@
using System;
using Elsa.HostedServices;
using Elsa.Services;
using Microsoft.Extensions.DependencyInjection;

View file

@ -14,7 +14,7 @@ namespace Elsa.Services
public async Task ScheduleTask(string channelName, Func<ValueTask> 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<Func<ValueTask>> GetOrCreateChannel(string channelName, Duration? channelTtl, CancellationToken cancellationToken)
private async Task<Channel<Func<ValueTask>>> GetOrCreateChannel(string channelName, Duration? channelTtl, CancellationToken cancellationToken)
{
_semaphore.Wait(cancellationToken);
await _semaphore.WaitAsync(cancellationToken);
try
{

View file

@ -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<string, IBus> _serviceBuses = new Dictionary<string, IBus>();
private readonly ConcurrentDictionary<string, IBus> _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;
}

View file

@ -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<EntityFrameworkWorkflowExecutionLogRecordStore>();
if (autorunMigrations)
elsa.Services.AddStartupTask<RunEFCoreMigrations>();
elsa.Services.AddHostedService<RunEFCoreMigrations>();
return elsa
.UseWorkflowDefinitionStore(sp => sp.GetRequiredService<EntityFrameworkWorkflowDefinitionStore>())

View file

@ -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
{
/// <summary>
/// Executes EF Core migrations.
/// </summary>
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<ElsaContext>();
await elsaContext.Database.MigrateAsync(cancellationToken);
}
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
}
}