Incremental work on service bus API

This commit is contained in:
Sipke Schoorstra 2020-11-20 17:08:05 +01:00
parent 431e2797b0
commit 355038f27e
11 changed files with 92 additions and 49 deletions

View file

@ -13,14 +13,14 @@ namespace Elsa.Activities.Rebus.Extensions
.AddActivity<SendMessage>() .AddActivity<SendMessage>()
.AddActivity<MessageReceived>(); .AddActivity<MessageReceived>();
public static IServiceCollection AddRebusActivities<T>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T>(); // public static IServiceCollection AddRebusActivities<T>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T>();
public static IServiceCollection AddRebusActivities<T1, T2>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>(); // public static IServiceCollection AddRebusActivities<T1, T2>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>();
public static IServiceCollection AddRebusActivities<T1, T2, T3>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>().AddMessageType<T3>(); // public static IServiceCollection AddRebusActivities<T1, T2, T3>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>().AddMessageType<T3>();
public static IServiceCollection AddRebusActivities<T1, T2, T3, T4>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>().AddMessageType<T3>().AddMessageType<T4>(); // public static IServiceCollection AddRebusActivities<T1, T2, T3, T4>(this IServiceCollection services) => services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>().AddMessageType<T3>().AddMessageType<T4>();
//
public static IServiceCollection AddRebusActivities<T1, T2, T3, T4, T5>(this IServiceCollection services) => // public static IServiceCollection AddRebusActivities<T1, T2, T3, T4, T5>(this IServiceCollection services) =>
services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>().AddMessageType<T3>().AddMessageType<T4>().AddMessageType<T5>(); // services.AddRebusActivities().AddMessageType<T1>().AddMessageType<T2>().AddMessageType<T3>().AddMessageType<T4>().AddMessageType<T5>();
//
public static IServiceCollection AddMessageType<T>(this IServiceCollection services) => services.AddConsumer<T, MessageConsumer<T>>(); // public static IServiceCollection AddMessageType<T>(this IServiceCollection services) => services.AddConsumer<T, MessageConsumer<T>>();
} }
} }

View file

@ -1,3 +1,4 @@
using System;
using System.Threading.Tasks; using System.Threading.Tasks;
using Elsa.Exceptions; using Elsa.Exceptions;
using Elsa.Extensions; using Elsa.Extensions;
@ -7,7 +8,7 @@ using Rebus.Handlers;
namespace Elsa.Consumers namespace Elsa.Consumers
{ {
public class RunWorkflowConsumer : IHandleMessages<RunWorkflow> public class RunWorkflowConsumer : IHandleMessages<RunWorkflow>, IDisposable
{ {
private readonly IWorkflowRunner _workflowRunner; private readonly IWorkflowRunner _workflowRunner;
private readonly IWorkflowInstanceManager _workflowInstanceManager; private readonly IWorkflowInstanceManager _workflowInstanceManager;
@ -30,5 +31,10 @@ namespace Elsa.Consumers
message.ActivityId, message.ActivityId,
message.Input); message.Input);
} }
public void Dispose()
{
}
} }
} }

View file

@ -3,6 +3,7 @@ using System.Data;
using Elsa.Caching; using Elsa.Caching;
using Elsa.DistributedLock; using Elsa.DistributedLock;
using Elsa.Extensions; using Elsa.Extensions;
using Elsa.Messages;
using Elsa.Serialization; using Elsa.Serialization;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Newtonsoft.Json; using Newtonsoft.Json;
@ -126,8 +127,8 @@ namespace Elsa
.Logging(logging => logging.ColoredConsole(LogLevel.Debug)) .Logging(logging => logging.ColoredConsole(LogLevel.Debug))
.Subscriptions(s => s.StoreInMemory(subscriberStore)) .Subscriptions(s => s.StoreInMemory(subscriberStore))
.DataBus(s => s.StoreInMemory(dataStore)) .DataBus(s => s.StoreInMemory(dataStore))
.Routing(r => r.TypeBased()) //.Routing(r => r.TypeBased().Map<RunWorkflow>("runworkflow"))
.Transport(t => t.UseInMemoryTransport(network, "elsa_publisher")); .Transport(t => t.UseInMemoryTransport(network, "elsa:publisher"));
} }
} }
} }

View file

@ -29,7 +29,9 @@ using NodaTime;
using Rebus.Bus; using Rebus.Bus;
using Rebus.Config; using Rebus.Config;
using Rebus.Handlers; using Rebus.Handlers;
using Rebus.Routing.TypeBased;
using Rebus.ServiceProvider; using Rebus.ServiceProvider;
using Rebus.Transport.InMem;
using RunWorkflow = Elsa.Messages.RunWorkflow; using RunWorkflow = Elsa.Messages.RunWorkflow;
// ReSharper disable once CheckNamespace // ReSharper disable once CheckNamespace
@ -81,26 +83,12 @@ namespace Microsoft.Extensions.DependencyInjection
.AddTransient(sp => workflow); .AddTransient(sp => workflow);
} }
public static IServiceCollection AddConsumer<TMessage, TConsumer>(this IServiceCollection services, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> configureBus, Action<IServiceCollection>? configureBusServices = default) public static void AddConsumer<TMessage, TConsumer>(this IServiceProvider sp, string name, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> configureBus)
where TConsumer : class, IHandleMessages<TMessage> where TConsumer : class, IHandleMessages<TMessage>
{ {
return services var factory = sp.GetRequiredService<IServiceBusContainerFactory>();
.AddTransient(sp => var container = factory.CreateBusContainer(name, configureBus);
{ container.Bus.Subscribe<TMessage>();
var factory = sp.GetRequiredService<IServiceBusContainerFactory>();
var container = factory.CreateBusContainer(containerServices =>
{
containerServices
.AddTransient<TConsumer>()
.AddTransient<IHandleMessages>(bsp => bsp.GetRequiredService<TConsumer>())
.AddTransient<IHandleMessages<TMessage>, TConsumer>(bsp => bsp.GetRequiredService<TConsumer>());
configureBusServices?.Invoke(containerServices);
containerServices.AddRebus(configureBus);
});
return container.ServiceProvider.GetRequiredService<TConsumer>();
});
} }
private static IServiceCollection AddMediatR(this ElsaOptions options) => options.Services.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity)); private static IServiceCollection AddMediatR(this ElsaOptions options) => options.Services.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity));
@ -148,7 +136,10 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton<ICloner, AutoMapperCloner>() .AddSingleton<ICloner, AutoMapperCloner>()
.AddNotificationHandlers(typeof(ElsaServiceCollectionExtensions)) .AddNotificationHandlers(typeof(ElsaServiceCollectionExtensions))
.AddStartupTask<StartServiceBusTask>() .AddStartupTask<StartServiceBusTask>()
.AddConsumer<RunWorkflow, RunWorkflowConsumer>() .AddSingleton<IServiceBusContainerFactory, ServiceBusContainerFactory>()
.AddTransient<RunWorkflowConsumer>()
.AddTransient<IHandleMessages>(bsp => bsp.GetRequiredService<RunWorkflowConsumer>())
.AddTransient<IHandleMessages<RunWorkflow>, RunWorkflowConsumer>(bsp => bsp.GetRequiredService<RunWorkflowConsumer>())
.AddMetadataHandlers() .AddMetadataHandlers()
.AddCoreActivities(); .AddCoreActivities();

View file

@ -5,7 +5,6 @@ namespace Elsa.ServiceBus
{ {
public interface IServiceBusContainer public interface IServiceBusContainer
{ {
IServiceProvider ServiceProvider { get; }
IBus Bus { get; } IBus Bus { get; }
} }
} }

View file

@ -1,10 +1,12 @@
using System; using System;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Rebus.Config;
namespace Elsa.ServiceBus namespace Elsa.ServiceBus
{ {
public interface IServiceBusContainerFactory public interface IServiceBusContainerFactory
{ {
IServiceBusContainer CreateBusContainer(Action<IServiceCollection> configure); IServiceBusContainer CreateBusContainer(string name, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> configure);
IServiceBusContainer GetBusContainer(string name);
} }
} }

View file

@ -1,22 +1,25 @@
using System; using System;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Rebus.Bus; using Rebus.Bus;
using Rebus.Config;
using Rebus.ServiceProvider;
namespace Elsa.ServiceBus namespace Elsa.ServiceBus
{ {
public class ServiceBusContainer : IServiceBusContainer, IDisposable public class ServiceBusContainer : IServiceBusContainer, IDisposable
{ {
public ServiceBusContainer(Action<IServiceCollection> configure) public ServiceBusContainer(string name, IServiceProvider serviceProvider, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> configure)
{ {
var services = new ServiceCollection(); Name = name;
configure(services);
ServiceProvider = services.BuildServiceProvider(); var configurer = Configure.With(new DependencyInjectionHandlerActivator(serviceProvider));
Bus = ServiceProvider.GetRequiredService<IBus>(); configurer = configure(configurer, serviceProvider);
Bus = configurer.Start();
} }
public IServiceProvider ServiceProvider { get; } public string Name { get; }
public IBus Bus { get; } public IBus Bus { get; }
public void Dispose() => ((IDisposable)ServiceProvider).Dispose(); public void Dispose() => Bus.Dispose();
} }
} }

View file

@ -1,18 +1,27 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Rebus.Config;
namespace Elsa.ServiceBus namespace Elsa.ServiceBus
{ {
public class ServiceBusContainerFactory : IServiceBusContainerFactory public class ServiceBusContainerFactory : IServiceBusContainerFactory
{ {
private readonly ICollection<IServiceBusContainer> _containers = new List<IServiceBusContainer>(); private readonly IServiceProvider _serviceProvider;
private readonly IDictionary<string, IServiceBusContainer> _containers = new Dictionary<string, IServiceBusContainer>();
public IServiceBusContainer CreateBusContainer(Action<IServiceCollection> configure) public ServiceBusContainerFactory(IServiceProvider serviceProvider)
{ {
var container = new ServiceBusContainer(configure); _serviceProvider = serviceProvider;
_containers.Add(container); }
public IServiceBusContainer CreateBusContainer(string name, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> configure)
{
var container = new ServiceBusContainer(name, _serviceProvider, configure);
_containers.Add(name, container);
return container; return container;
} }
public IServiceBusContainer GetBusContainer(string name) => _containers[name];
} }
} }

View file

@ -10,6 +10,7 @@ using Elsa.Extensions;
using Elsa.Indexes; using Elsa.Indexes;
using Elsa.Messages; using Elsa.Messages;
using Elsa.Models; using Elsa.Models;
using Elsa.ServiceBus;
using Elsa.Services.Models; using Elsa.Services.Models;
using Elsa.Triggers; using Elsa.Triggers;
using MediatR; using MediatR;
@ -26,6 +27,7 @@ namespace Elsa.Services
private readonly IWorkflowSelector _workflowSelector; private readonly IWorkflowSelector _workflowSelector;
private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowRegistry _workflowRegistry;
private readonly IWorkflowSchedulerQueue _queue; private readonly IWorkflowSchedulerQueue _queue;
private readonly IServiceBusContainerFactory _serviceBusContainerFactory;
public WorkflowScheduler( public WorkflowScheduler(
IBus serviceBus, IBus serviceBus,
@ -33,7 +35,8 @@ namespace Elsa.Services
IWorkflowFactory workflowFactory, IWorkflowFactory workflowFactory,
IWorkflowSelector workflowSelector, IWorkflowSelector workflowSelector,
IWorkflowRegistry workflowRegistry, IWorkflowRegistry workflowRegistry,
IWorkflowSchedulerQueue queue) IWorkflowSchedulerQueue queue,
IServiceBusContainerFactory serviceBusContainerFactory)
{ {
_serviceBus = serviceBus; _serviceBus = serviceBus;
_workflowInstanceManager = workflowInstanceManager; _workflowInstanceManager = workflowInstanceManager;
@ -41,14 +44,18 @@ namespace Elsa.Services
_workflowSelector = workflowSelector; _workflowSelector = workflowSelector;
_workflowRegistry = workflowRegistry; _workflowRegistry = workflowRegistry;
_queue = queue; _queue = queue;
_serviceBusContainerFactory = serviceBusContainerFactory;
} }
public async Task ScheduleWorkflowInstanceAsync( public async Task ScheduleWorkflowInstanceAsync(
string instanceId, string instanceId,
string? activityId = default, string? activityId = default,
object? input = default, object? input = default,
CancellationToken cancellationToken = default) => CancellationToken cancellationToken = default)
await _serviceBus.Send(new RunWorkflow(instanceId, activityId, input)); {
var bus = _serviceBusContainerFactory.GetBusContainer("run_workflow:client");
await bus.Bus.Send(new RunWorkflow(instanceId, activityId, input));
}
public async Task<IEnumerable<WorkflowInstance>> ScheduleWorkflowDefinitionAsync( public async Task<IEnumerable<WorkflowInstance>> ScheduleWorkflowDefinitionAsync(
string definitionId, string definitionId,

View file

@ -2,10 +2,16 @@ using System;
using System.Linq; using System.Linq;
using System.Threading; using System.Threading;
using System.Threading.Tasks; using System.Threading.Tasks;
using Elsa.Consumers;
using Elsa.Messages;
using Elsa.Services; using Elsa.Services;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Rebus.DataBus.InMem;
using Rebus.Handlers; using Rebus.Handlers;
using Rebus.Persistence.InMem;
using Rebus.Routing.TypeBased;
using Rebus.ServiceProvider; using Rebus.ServiceProvider;
using Rebus.Transport.InMem;
namespace Elsa.StartupTasks namespace Elsa.StartupTasks
{ {
@ -25,6 +31,24 @@ namespace Elsa.StartupTasks
await bus.Subscribe(messageType); await bus.Subscribe(messageType);
}); });
var transport = _serviceProvider.GetService<InMemNetwork>();
var store = _serviceProvider.GetRequiredService<InMemorySubscriberStore>();
var queueName = "run_Workflow";
// Sender
_serviceProvider
.AddConsumer<RunWorkflow, RunWorkflowConsumer>("run_workflow:sender", (bus, _) => bus
.Logging(l => l.ColoredConsole())
.Transport(t => t.UseInMemoryTransportAsOneWayClient(transport)));
// Receiver
_serviceProvider
.AddConsumer<RunWorkflow, RunWorkflowConsumer>("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<RunWorkflow>(queueName)));
return Task.CompletedTask; return Task.CompletedTask;
} }
} }

View file

@ -24,7 +24,8 @@ namespace Elsa.Samples.RebusWorker
.UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared"))) .UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared")))
.AddConsoleActivities() .AddConsoleActivities()
.AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1)) .AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1))
.AddRebusActivities<Greeting>() .AddRebusActivities()
//.AddConsumer<Greeting, Gree>()
.AddWorkflow<ProducerWorkflow>() .AddWorkflow<ProducerWorkflow>()
.AddWorkflow<ConsumerWorkflow>(); .AddWorkflow<ConsumerWorkflow>();
}); });