diff --git a/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs index 9e48e40ff..43ed51664 100644 --- a/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs @@ -13,14 +13,14 @@ namespace Elsa.Activities.Rebus.Extensions .AddActivity() .AddActivity(); - public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType(); - public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType(); - public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType(); - public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType().AddMessageType(); - - public static IServiceCollection AddRebusActivities(this IServiceCollection services) => - services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType().AddMessageType().AddMessageType(); - - public static IServiceCollection AddMessageType(this IServiceCollection services) => services.AddConsumer>(); + // public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType(); + // public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType(); + // public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType(); + // public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType().AddMessageType(); + // + // public static IServiceCollection AddRebusActivities(this IServiceCollection services) => + // services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType().AddMessageType().AddMessageType(); + // + // public static IServiceCollection AddMessageType(this IServiceCollection services) => services.AddConsumer>(); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowConsumer.cs index 29e537044..c42d9e175 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowConsumer.cs @@ -1,3 +1,4 @@ +using System; using System.Threading.Tasks; using Elsa.Exceptions; using Elsa.Extensions; @@ -7,7 +8,7 @@ using Rebus.Handlers; namespace Elsa.Consumers { - public class RunWorkflowConsumer : IHandleMessages + public class RunWorkflowConsumer : IHandleMessages, IDisposable { private readonly IWorkflowRunner _workflowRunner; private readonly IWorkflowInstanceManager _workflowInstanceManager; @@ -30,5 +31,10 @@ namespace Elsa.Consumers message.ActivityId, message.Input); } + + public void Dispose() + { + + } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/ElsaOptions.cs b/src/core/Elsa.Core/ElsaOptions.cs index 886407361..110ac205e 100644 --- a/src/core/Elsa.Core/ElsaOptions.cs +++ b/src/core/Elsa.Core/ElsaOptions.cs @@ -3,6 +3,7 @@ using System.Data; using Elsa.Caching; using Elsa.DistributedLock; using Elsa.Extensions; +using Elsa.Messages; using Elsa.Serialization; using Microsoft.Extensions.DependencyInjection; using Newtonsoft.Json; @@ -126,8 +127,8 @@ namespace Elsa .Logging(logging => logging.ColoredConsole(LogLevel.Debug)) .Subscriptions(s => s.StoreInMemory(subscriberStore)) .DataBus(s => s.StoreInMemory(dataStore)) - .Routing(r => r.TypeBased()) - .Transport(t => t.UseInMemoryTransport(network, "elsa_publisher")); + //.Routing(r => r.TypeBased().Map("runworkflow")) + .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 cfb6ea86e..230bcad51 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -29,7 +29,9 @@ using NodaTime; using Rebus.Bus; using Rebus.Config; using Rebus.Handlers; +using Rebus.Routing.TypeBased; using Rebus.ServiceProvider; +using Rebus.Transport.InMem; using RunWorkflow = Elsa.Messages.RunWorkflow; // ReSharper disable once CheckNamespace @@ -81,26 +83,12 @@ namespace Microsoft.Extensions.DependencyInjection .AddTransient(sp => workflow); } - public static IServiceCollection AddConsumer(this IServiceCollection services, Func configureBus, Action? configureBusServices = default) + public static void AddConsumer(this IServiceProvider sp, string name, Func configureBus) where TConsumer : class, IHandleMessages { - return services - .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(); - }); + var factory = sp.GetRequiredService(); + var container = factory.CreateBusContainer(name, configureBus); + container.Bus.Subscribe(); } private static IServiceCollection AddMediatR(this ElsaOptions options) => options.Services.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity)); @@ -148,7 +136,10 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton() .AddNotificationHandlers(typeof(ElsaServiceCollectionExtensions)) .AddStartupTask() - .AddConsumer() + .AddSingleton() + .AddTransient() + .AddTransient(bsp => bsp.GetRequiredService()) + .AddTransient, RunWorkflowConsumer>(bsp => bsp.GetRequiredService()) .AddMetadataHandlers() .AddCoreActivities(); diff --git a/src/core/Elsa.Core/ServiceBus/IServiceBusContainer.cs b/src/core/Elsa.Core/ServiceBus/IServiceBusContainer.cs index ae2be6763..8af07b7cb 100644 --- a/src/core/Elsa.Core/ServiceBus/IServiceBusContainer.cs +++ b/src/core/Elsa.Core/ServiceBus/IServiceBusContainer.cs @@ -5,7 +5,6 @@ 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 index 4f7bbdceb..8147d58f4 100644 --- a/src/core/Elsa.Core/ServiceBus/IServiceBusContainerFactory.cs +++ b/src/core/Elsa.Core/ServiceBus/IServiceBusContainerFactory.cs @@ -1,10 +1,12 @@ using System; using Microsoft.Extensions.DependencyInjection; +using Rebus.Config; namespace Elsa.ServiceBus { public interface IServiceBusContainerFactory { - IServiceBusContainer CreateBusContainer(Action configure); + IServiceBusContainer CreateBusContainer(string name, Func configure); + IServiceBusContainer GetBusContainer(string name); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/ServiceBus/ServiceBusContainer.cs b/src/core/Elsa.Core/ServiceBus/ServiceBusContainer.cs index 76291b334..b417ccb64 100644 --- a/src/core/Elsa.Core/ServiceBus/ServiceBusContainer.cs +++ b/src/core/Elsa.Core/ServiceBus/ServiceBusContainer.cs @@ -1,22 +1,25 @@ using System; using Microsoft.Extensions.DependencyInjection; using Rebus.Bus; +using Rebus.Config; +using Rebus.ServiceProvider; namespace Elsa.ServiceBus { public class ServiceBusContainer : IServiceBusContainer, IDisposable { - public ServiceBusContainer(Action configure) + public ServiceBusContainer(string name, IServiceProvider serviceProvider, Func configure) { - var services = new ServiceCollection(); - configure(services); - ServiceProvider = services.BuildServiceProvider(); - Bus = ServiceProvider.GetRequiredService(); + Name = name; + + var configurer = Configure.With(new DependencyInjectionHandlerActivator(serviceProvider)); + configurer = configure(configurer, serviceProvider); + Bus = configurer.Start(); } - public IServiceProvider ServiceProvider { get; } + public string Name { get; } public IBus Bus { get; } - public void Dispose() => ((IDisposable)ServiceProvider).Dispose(); + public void Dispose() => Bus.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 index 8d195cd46..ffd76f676 100644 --- a/src/core/Elsa.Core/ServiceBus/ServiceBusContainerFactory.cs +++ b/src/core/Elsa.Core/ServiceBus/ServiceBusContainerFactory.cs @@ -1,18 +1,27 @@ using System; using System.Collections.Generic; using Microsoft.Extensions.DependencyInjection; +using Rebus.Config; namespace Elsa.ServiceBus { - public class ServiceBusContainerFactory : IServiceBusContainerFactory + public class ServiceBusContainerFactory : IServiceBusContainerFactory { - private readonly ICollection _containers = new List(); + private readonly IServiceProvider _serviceProvider; + private readonly IDictionary _containers = new Dictionary(); - public IServiceBusContainer CreateBusContainer(Action configure) + public ServiceBusContainerFactory(IServiceProvider serviceProvider) { - var container = new ServiceBusContainer(configure); - _containers.Add(container); + _serviceProvider = serviceProvider; + } + + public IServiceBusContainer CreateBusContainer(string name, Func configure) + { + var container = new ServiceBusContainer(name, _serviceProvider, configure); + _containers.Add(name, container); return container; } + + public IServiceBusContainer GetBusContainer(string name) => _containers[name]; } } \ 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 36165bf75..e48603233 100644 --- a/src/core/Elsa.Core/Services/WorkflowScheduler.cs +++ b/src/core/Elsa.Core/Services/WorkflowScheduler.cs @@ -10,6 +10,7 @@ using Elsa.Extensions; using Elsa.Indexes; using Elsa.Messages; using Elsa.Models; +using Elsa.ServiceBus; using Elsa.Services.Models; using Elsa.Triggers; using MediatR; @@ -26,6 +27,7 @@ namespace Elsa.Services private readonly IWorkflowSelector _workflowSelector; private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowSchedulerQueue _queue; + private readonly IServiceBusContainerFactory _serviceBusContainerFactory; public WorkflowScheduler( IBus serviceBus, @@ -33,7 +35,8 @@ namespace Elsa.Services IWorkflowFactory workflowFactory, IWorkflowSelector workflowSelector, IWorkflowRegistry workflowRegistry, - IWorkflowSchedulerQueue queue) + IWorkflowSchedulerQueue queue, + IServiceBusContainerFactory serviceBusContainerFactory) { _serviceBus = serviceBus; _workflowInstanceManager = workflowInstanceManager; @@ -41,14 +44,18 @@ namespace Elsa.Services _workflowSelector = workflowSelector; _workflowRegistry = workflowRegistry; _queue = queue; + _serviceBusContainerFactory = serviceBusContainerFactory; } public async Task ScheduleWorkflowInstanceAsync( string instanceId, string? activityId = default, object? input = default, - CancellationToken cancellationToken = default) => - await _serviceBus.Send(new RunWorkflow(instanceId, activityId, input)); + CancellationToken cancellationToken = default) + { + var bus = _serviceBusContainerFactory.GetBusContainer("run_workflow:client"); + await bus.Bus.Send(new RunWorkflow(instanceId, activityId, input)); + } public async Task> ScheduleWorkflowDefinitionAsync( string definitionId, diff --git a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs index d6d67b115..c7569c72d 100644 --- a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs +++ b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs @@ -2,10 +2,16 @@ using System; using System.Linq; using System.Threading; using System.Threading.Tasks; +using Elsa.Consumers; +using Elsa.Messages; using Elsa.Services; using Microsoft.Extensions.DependencyInjection; +using Rebus.DataBus.InMem; using Rebus.Handlers; +using Rebus.Persistence.InMem; +using Rebus.Routing.TypeBased; using Rebus.ServiceProvider; +using Rebus.Transport.InMem; namespace Elsa.StartupTasks { @@ -25,6 +31,24 @@ namespace Elsa.StartupTasks await bus.Subscribe(messageType); }); + var transport = _serviceProvider.GetService(); + var store = _serviceProvider.GetRequiredService(); + var queueName = "run_Workflow"; + + // Sender + _serviceProvider + .AddConsumer("run_workflow:sender", (bus, _) => bus + .Logging(l => l.ColoredConsole()) + .Transport(t => t.UseInMemoryTransportAsOneWayClient(transport))); + + // Receiver + _serviceProvider + .AddConsumer("run_workflow:client", (bus, _) => bus + .Logging(l => l.ColoredConsole()) + .Subscriptions(s => s.StoreInMemory(store)) + .Transport(t => t.UseInMemoryTransport(transport, queueName)) + .Routing(r => r.TypeBased().Map(queueName))); + return Task.CompletedTask; } } diff --git a/src/samples/Elsa.Samples.RebusWorker/Program.cs b/src/samples/Elsa.Samples.RebusWorker/Program.cs index 43a73ff29..6e8534147 100644 --- a/src/samples/Elsa.Samples.RebusWorker/Program.cs +++ b/src/samples/Elsa.Samples.RebusWorker/Program.cs @@ -24,7 +24,8 @@ namespace Elsa.Samples.RebusWorker .UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared"))) .AddConsoleActivities() .AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1)) - .AddRebusActivities() + .AddRebusActivities() + //.AddConsumer() .AddWorkflow() .AddWorkflow(); });