diff --git a/src/core/Elsa.Core/ElsaOptions.cs b/src/core/Elsa.Core/ElsaOptions.cs index ba1ec407d..886407361 100644 --- a/src/core/Elsa.Core/ElsaOptions.cs +++ b/src/core/Elsa.Core/ElsaOptions.cs @@ -11,6 +11,7 @@ using Rebus.DataBus.InMem; using Rebus.Logging; using Rebus.Persistence.InMem; using Rebus.Routing.TypeBased; +using Rebus.ServiceProvider; using Rebus.Transport.InMem; using YesSql; using Storage.Net; @@ -29,7 +30,6 @@ namespace Elsa StorageFactory = sp => Storage.Net.StorageFactory.Blobs.InMemory(); DistributedLockProviderFactory = sp => new DefaultLockProvider(); SignalFactory = sp => new Signal(); - ServiceBusConfigurer = ConfigureInMemoryServiceBus; JsonSerializerConfigurer = (sp, serializer) => {}; AddAutoMapper = () => @@ -37,6 +37,14 @@ namespace Elsa services.AddAutoMapper(ServiceLifetime.Singleton); services.AddSingleton(sp => sp.CreateAutoMapperConfiguration()); }; + + AddServiceBus = () => + { + services.AddSingleton(); + services.AddSingleton(); + services.AddSingleton(); + services.AddRebus(ConfigureInMemoryServiceBus); + }; CreateJsonSerializer = sp => { @@ -51,10 +59,10 @@ namespace Elsa internal Func StorageFactory { get; set; } internal Func DistributedLockProviderFactory { get; private set; } internal Func SignalFactory { get; private set; } - internal Func ServiceBusConfigurer { get; private set; } internal Func CreateJsonSerializer { get; private set; } internal Action JsonSerializerConfigurer { get; private set; } internal Action AddAutoMapper { get; private set; } + internal Action AddServiceBus { get; private set; } public ElsaOptions UseDistributedLockProvider(Func factory) { @@ -90,14 +98,12 @@ namespace Elsa return this; } - public ElsaOptions ConfigureServiceBus(Func configure) + public ElsaOptions UseServiceBus(Action addServiceBus) { - ServiceBusConfigurer = configure; + AddServiceBus = addServiceBus; return this; } - public ElsaOptions ConfigureServiceBus(Func configure) => ConfigureServiceBus((bus, _) => configure(bus)); - public ElsaOptions UseJsonSerializer(Func factory) { CreateJsonSerializer = factory; @@ -112,12 +118,16 @@ namespace Elsa private static RebusConfigurer ConfigureInMemoryServiceBus(RebusConfigurer rebus, IServiceProvider serviceProvider) { + var subscriberStore = serviceProvider.GetRequiredService(); + var dataStore = serviceProvider.GetRequiredService(); + var network = serviceProvider.GetRequiredService(); + return rebus - .Logging(logging => logging.ColoredConsole(LogLevel.Info)) - .Subscriptions(s => s.StoreInMemory(new InMemorySubscriberStore())) - .DataBus(s => s.StoreInMemory(new InMemDataStore())) + .Logging(logging => logging.ColoredConsole(LogLevel.Debug)) + .Subscriptions(s => s.StoreInMemory(subscriberStore)) + .DataBus(s => s.StoreInMemory(dataStore)) .Routing(r => r.TypeBased()) - .Transport(t => t.UseInMemoryTransport(new InMemNetwork(), "Messages")); + .Transport(t => t.UseInMemoryTransport(network, "elsa_publisher")); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index f1c361319..cfb6ea86e 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -26,7 +26,10 @@ using Elsa.WorkflowProviders; using MediatR; using Microsoft.Extensions.DependencyInjection.Extensions; using NodaTime; +using Rebus.Bus; +using Rebus.Config; using Rebus.Handlers; +using Rebus.ServiceProvider; using RunWorkflow = Elsa.Messages.RunWorkflow; // ReSharper disable once CheckNamespace @@ -78,12 +81,26 @@ namespace Microsoft.Extensions.DependencyInjection .AddTransient(sp => workflow); } - public static IServiceCollection AddConsumer(this IServiceCollection services) where TConsumer : class, IHandleMessages + public static IServiceCollection AddConsumer(this IServiceCollection services, Func configureBus, Action? configureBusServices = default) + where TConsumer : class, IHandleMessages { return services - .AddTransient() - .AddTransient(sp => sp.GetRequiredService()) - .AddTransient, TConsumer>(sp => sp.GetRequiredService()); + .AddTransient(sp => + { + var factory = sp.GetRequiredService(); + var container = factory.CreateBusContainer(containerServices => + { + containerServices + .AddTransient() + .AddTransient(bsp => bsp.GetRequiredService()) + .AddTransient, TConsumer>(bsp => bsp.GetRequiredService()); + + configureBusServices?.Invoke(containerServices); + containerServices.AddRebus(configureBus); + }); + + return container.ServiceProvider.GetRequiredService(); + }); } private static IServiceCollection AddMediatR(this ElsaOptions options) => options.Services.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity)); @@ -163,11 +180,5 @@ namespace Microsoft.Extensions.DependencyInjection .AddScoped() .AddActivity() .AddTriggerProvider(); - - private static ElsaOptions AddServiceBus(this ElsaOptions options) - { - options.WithServiceBus(options.ServiceBusConfigurer); - return options; - } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/ServiceBus/IServiceBusContainer.cs b/src/core/Elsa.Core/ServiceBus/IServiceBusContainer.cs new file mode 100644 index 000000000..ae2be6763 --- /dev/null +++ b/src/core/Elsa.Core/ServiceBus/IServiceBusContainer.cs @@ -0,0 +1,11 @@ +using System; +using Rebus.Bus; + +namespace Elsa.ServiceBus +{ + public interface IServiceBusContainer + { + IServiceProvider ServiceProvider { get; } + IBus Bus { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/ServiceBus/IServiceBusContainerFactory.cs b/src/core/Elsa.Core/ServiceBus/IServiceBusContainerFactory.cs new file mode 100644 index 000000000..4f7bbdceb --- /dev/null +++ b/src/core/Elsa.Core/ServiceBus/IServiceBusContainerFactory.cs @@ -0,0 +1,10 @@ +using System; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.ServiceBus +{ + public interface IServiceBusContainerFactory + { + IServiceBusContainer CreateBusContainer(Action configure); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/ServiceBus/ServiceBusContainer.cs b/src/core/Elsa.Core/ServiceBus/ServiceBusContainer.cs new file mode 100644 index 000000000..76291b334 --- /dev/null +++ b/src/core/Elsa.Core/ServiceBus/ServiceBusContainer.cs @@ -0,0 +1,22 @@ +using System; +using Microsoft.Extensions.DependencyInjection; +using Rebus.Bus; + +namespace Elsa.ServiceBus +{ + public class ServiceBusContainer : IServiceBusContainer, IDisposable + { + public ServiceBusContainer(Action configure) + { + var services = new ServiceCollection(); + configure(services); + ServiceProvider = services.BuildServiceProvider(); + Bus = ServiceProvider.GetRequiredService(); + } + + public IServiceProvider ServiceProvider { get; } + public IBus Bus { get; } + + public void Dispose() => ((IDisposable)ServiceProvider).Dispose(); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/ServiceBus/ServiceBusContainerFactory.cs b/src/core/Elsa.Core/ServiceBus/ServiceBusContainerFactory.cs new file mode 100644 index 000000000..8d195cd46 --- /dev/null +++ b/src/core/Elsa.Core/ServiceBus/ServiceBusContainerFactory.cs @@ -0,0 +1,18 @@ +using System; +using System.Collections.Generic; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.ServiceBus +{ + public class ServiceBusContainerFactory : IServiceBusContainerFactory + { + private readonly ICollection _containers = new List(); + + public IServiceBusContainer CreateBusContainer(Action configure) + { + var container = new ServiceBusContainer(configure); + _containers.Add(container); + return container; + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/ServiceBus/ServiceCollectionExtensions.cs b/src/core/Elsa.Core/ServiceBus/ServiceCollectionExtensions.cs index 3e7b3badc..3c8e4b6e1 100644 --- a/src/core/Elsa.Core/ServiceBus/ServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/ServiceBus/ServiceCollectionExtensions.cs @@ -1,5 +1,6 @@ using System; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; using Rebus.Config; using Rebus.ServiceProvider; @@ -7,7 +8,7 @@ namespace Elsa.ServiceBus { public static class ServiceCollectionExtensions { - public static IServiceCollection WithServiceBus(this ElsaOptions configuration, Func setup) => configuration.Services.AddRebus(setup); - public static IServiceCollection WithServiceBus(this ElsaOptions configuration, Func setup) => configuration.Services.AddRebus(setup); + // public static IServiceCollection WithServiceBus(this ElsaOptions configuration, Func setup) => configuration.Services.AddRebus(setup); + // public static IServiceCollection WithServiceBus(this ElsaOptions configuration, Func setup) => configuration.Services.AddRebus(setup); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowScheduler.cs b/src/core/Elsa.Core/Services/WorkflowScheduler.cs index b489d7a89..36165bf75 100644 --- a/src/core/Elsa.Core/Services/WorkflowScheduler.cs +++ b/src/core/Elsa.Core/Services/WorkflowScheduler.cs @@ -48,7 +48,7 @@ namespace Elsa.Services string? activityId = default, object? input = default, CancellationToken cancellationToken = default) => - await _serviceBus.Publish(new RunWorkflow(instanceId, activityId, input)); + await _serviceBus.Send(new RunWorkflow(instanceId, activityId, input)); public async Task> ScheduleWorkflowDefinitionAsync( string definitionId, diff --git a/src/samples/Elsa.Samples.RebusWorker/Program.cs b/src/samples/Elsa.Samples.RebusWorker/Program.cs index 0af51fc43..43a73ff29 100644 --- a/src/samples/Elsa.Samples.RebusWorker/Program.cs +++ b/src/samples/Elsa.Samples.RebusWorker/Program.cs @@ -1,16 +1,9 @@ -using System; using Elsa.Activities.Rebus.Extensions; using Elsa.Samples.RebusWorker.Messages; using Elsa.Samples.RebusWorker.Workflows; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using NodaTime; -using Rebus.Config; -using Rebus.DataBus.InMem; -using Rebus.Logging; -using Rebus.Persistence.InMem; -using Rebus.Routing.TypeBased; -using Rebus.Transport.InMem; using YesSql.Provider.Sqlite; namespace Elsa.Samples.RebusWorker @@ -28,23 +21,12 @@ namespace Elsa.Samples.RebusWorker { services .AddElsa(option => option - .UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared")) - .ConfigureServiceBus(ConfigureRebus)) + .UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared"))) .AddConsoleActivities() .AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1)) .AddRebusActivities() .AddWorkflow() .AddWorkflow(); }); - - private static RebusConfigurer ConfigureRebus(RebusConfigurer rebus, IServiceProvider serviceProvider) - { - return rebus - .Logging(logging => logging.ColoredConsole(LogLevel.Info)) - .Subscriptions(s => s.StoreInMemory(new InMemorySubscriberStore())) - .DataBus(s => s.StoreInMemory(new InMemDataStore())) - .Routing(r => r.TypeBased().Map("greeting")) - .Transport(t => t.UseInMemoryTransport(new InMemNetwork(), "inbox")); - } } } \ No newline at end of file diff --git a/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs b/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs index c8bda1957..b5132448c 100644 --- a/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs +++ b/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs @@ -1,4 +1,5 @@ using System; +using System.Linq; using Elsa.Activities.Console; using Elsa.Activities.Rebus; using Elsa.Activities.Timers; @@ -18,7 +19,7 @@ namespace Elsa.Samples.RebusWorker.Workflows _clock = clock; _random = new Random(); } - + public void Build(IWorkflowBuilder workflow) { workflow @@ -31,30 +32,18 @@ namespace Elsa.Samples.RebusWorker.Workflows private Greeting GetRandomGreeting() { - var greetings = new[] - { - new Greeting - { - From = "John", - To = "Jill", - Message = "Hello!" - }, - new Greeting - { - From = "Julia", - To = "Miriam", - Message = "Happy Monday!" - }, - new Greeting - { - From = "Jack", - To = "Bob", - Message = "How do you do?" - } - }; + var names = new[] { "John", "Jill", "Julia", "Miriam", "Jack", "Bob" }; + var messages = new[] { "Hello!", "How do you do?", "Happy Monday!" }; + var from = _random.Next(0, names.Length); + var to = _random.Next(0, names.Length); + var message = _random.Next(0, messages.Length); - var index = _random.Next(0, greetings.Length); - return greetings[index]; + return new Greeting + { + From = names[from], + To = names[to], + Message = messages[message] + }; } } } \ No newline at end of file