This commit is contained in:
Sipke Schoorstra 2020-11-20 13:29:08 +01:00
parent 2d54b7e337
commit d5633d282c
10 changed files with 120 additions and 66 deletions

View file

@ -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<InMemNetwork>();
services.AddSingleton<InMemorySubscriberStore>();
services.AddSingleton<InMemDataStore>();
services.AddRebus(ConfigureInMemoryServiceBus);
};
CreateJsonSerializer = sp =>
{
@ -51,10 +59,10 @@ namespace Elsa
internal Func<IServiceProvider, IBlobStorage> StorageFactory { get; set; }
internal Func<IServiceProvider, IDistributedLockProvider> DistributedLockProviderFactory { get; private set; }
internal Func<IServiceProvider, ISignal> SignalFactory { get; private set; }
internal Func<RebusConfigurer, IServiceProvider, RebusConfigurer> ServiceBusConfigurer { get; private set; }
internal Func<IServiceProvider, JsonSerializer> CreateJsonSerializer { get; private set; }
internal Action<IServiceProvider, JsonSerializer> JsonSerializerConfigurer { get; private set; }
internal Action AddAutoMapper { get; private set; }
internal Action AddServiceBus { get; private set; }
public ElsaOptions UseDistributedLockProvider(Func<IServiceProvider, IDistributedLockProvider> factory)
{
@ -90,14 +98,12 @@ namespace Elsa
return this;
}
public ElsaOptions ConfigureServiceBus(Func<RebusConfigurer, IServiceProvider, RebusConfigurer> configure)
public ElsaOptions UseServiceBus(Action addServiceBus)
{
ServiceBusConfigurer = configure;
AddServiceBus = addServiceBus;
return this;
}
public ElsaOptions ConfigureServiceBus(Func<RebusConfigurer, RebusConfigurer> configure) => ConfigureServiceBus((bus, _) => configure(bus));
public ElsaOptions UseJsonSerializer(Func<IServiceProvider, JsonSerializer> factory)
{
CreateJsonSerializer = factory;
@ -112,12 +118,16 @@ namespace Elsa
private static RebusConfigurer ConfigureInMemoryServiceBus(RebusConfigurer rebus, IServiceProvider serviceProvider)
{
var subscriberStore = serviceProvider.GetRequiredService<InMemorySubscriberStore>();
var dataStore = serviceProvider.GetRequiredService<InMemDataStore>();
var network = serviceProvider.GetRequiredService<InMemNetwork>();
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"));
}
}
}

View file

@ -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<TMessage, TConsumer>(this IServiceCollection services) where TConsumer : class, IHandleMessages<TMessage>
public static IServiceCollection AddConsumer<TMessage, TConsumer>(this IServiceCollection services, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> configureBus, Action<IServiceCollection>? configureBusServices = default)
where TConsumer : class, IHandleMessages<TMessage>
{
return services
.AddTransient<TConsumer>()
.AddTransient<IHandleMessages>(sp => sp.GetRequiredService<TConsumer>())
.AddTransient<IHandleMessages<TMessage>, TConsumer>(sp => sp.GetRequiredService<TConsumer>());
.AddTransient(sp =>
{
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));
@ -163,11 +180,5 @@ namespace Microsoft.Extensions.DependencyInjection
.AddScoped<ISignaler, Signaler>()
.AddActivity<Elsa.Activities.Workflows.RunWorkflow>()
.AddTriggerProvider<RunWorkflowTriggerProvider>();
private static ElsaOptions AddServiceBus(this ElsaOptions options)
{
options.WithServiceBus(options.ServiceBusConfigurer);
return options;
}
}
}

View file

@ -0,0 +1,11 @@
using System;
using Rebus.Bus;
namespace Elsa.ServiceBus
{
public interface IServiceBusContainer
{
IServiceProvider ServiceProvider { get; }
IBus Bus { get; }
}
}

View file

@ -0,0 +1,10 @@
using System;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.ServiceBus
{
public interface IServiceBusContainerFactory
{
IServiceBusContainer CreateBusContainer(Action<IServiceCollection> configure);
}
}

View file

@ -0,0 +1,22 @@
using System;
using Microsoft.Extensions.DependencyInjection;
using Rebus.Bus;
namespace Elsa.ServiceBus
{
public class ServiceBusContainer : IServiceBusContainer, IDisposable
{
public ServiceBusContainer(Action<IServiceCollection> configure)
{
var services = new ServiceCollection();
configure(services);
ServiceProvider = services.BuildServiceProvider();
Bus = ServiceProvider.GetRequiredService<IBus>();
}
public IServiceProvider ServiceProvider { get; }
public IBus Bus { get; }
public void Dispose() => ((IDisposable)ServiceProvider).Dispose();
}
}

View file

@ -0,0 +1,18 @@
using System;
using System.Collections.Generic;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.ServiceBus
{
public class ServiceBusContainerFactory : IServiceBusContainerFactory
{
private readonly ICollection<IServiceBusContainer> _containers = new List<IServiceBusContainer>();
public IServiceBusContainer CreateBusContainer(Action<IServiceCollection> configure)
{
var container = new ServiceBusContainer(configure);
_containers.Add(container);
return container;
}
}
}

View file

@ -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<RebusConfigurer, RebusConfigurer> setup) => configuration.Services.AddRebus(setup);
public static IServiceCollection WithServiceBus(this ElsaOptions configuration, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> setup) => configuration.Services.AddRebus(setup);
// public static IServiceCollection WithServiceBus(this ElsaOptions configuration, Func<RebusConfigurer, RebusConfigurer> setup) => configuration.Services.AddRebus(setup);
// public static IServiceCollection WithServiceBus(this ElsaOptions configuration, Func<RebusConfigurer, IServiceProvider, RebusConfigurer> setup) => configuration.Services.AddRebus(setup);
}
}

View file

@ -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<IEnumerable<WorkflowInstance>> ScheduleWorkflowDefinitionAsync(
string definitionId,

View file

@ -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<Greeting>()
.AddWorkflow<ProducerWorkflow>()
.AddWorkflow<ConsumerWorkflow>();
});
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>("greeting"))
.Transport(t => t.UseInMemoryTransport(new InMemNetwork(), "inbox"));
}
}
}

View file

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