Add Rebus activities

This commit is contained in:
Sipke Schoorstra 2020-11-19 22:21:17 +01:00
parent c586d070d8
commit 2d54b7e337
31 changed files with 479 additions and 38 deletions

View file

@ -138,6 +138,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "ElsaDashboard.Application.W
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.ForkJoinTimerAndSignalHttp", "src\samples\Elsa.Samples.ForkJoinTimerAndSignalHttp\Elsa.Samples.ForkJoinTimerAndSignalHttp.csproj", "{922F1EB6-5C8F-45DC-82A8-C651E4554542}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Activities.Rebus", "src\activities\Elsa.Activities.Rebus\Elsa.Activities.Rebus.csproj", "{115DEB38-679F-467E-8F87-9E669DFD2EB8}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.RebusWorker", "src\samples\Elsa.Samples.RebusWorker\Elsa.Samples.RebusWorker.csproj", "{A855C6B9-1548-4183-926C-75D80CEEF510}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -337,6 +341,14 @@ Global
{922F1EB6-5C8F-45DC-82A8-C651E4554542}.Debug|Any CPU.Build.0 = Debug|Any CPU
{922F1EB6-5C8F-45DC-82A8-C651E4554542}.Release|Any CPU.ActiveCfg = Release|Any CPU
{922F1EB6-5C8F-45DC-82A8-C651E4554542}.Release|Any CPU.Build.0 = Release|Any CPU
{115DEB38-679F-467E-8F87-9E669DFD2EB8}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{115DEB38-679F-467E-8F87-9E669DFD2EB8}.Debug|Any CPU.Build.0 = Debug|Any CPU
{115DEB38-679F-467E-8F87-9E669DFD2EB8}.Release|Any CPU.ActiveCfg = Release|Any CPU
{115DEB38-679F-467E-8F87-9E669DFD2EB8}.Release|Any CPU.Build.0 = Release|Any CPU
{A855C6B9-1548-4183-926C-75D80CEEF510}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{A855C6B9-1548-4183-926C-75D80CEEF510}.Debug|Any CPU.Build.0 = Debug|Any CPU
{A855C6B9-1548-4183-926C-75D80CEEF510}.Release|Any CPU.ActiveCfg = Release|Any CPU
{A855C6B9-1548-4183-926C-75D80CEEF510}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -402,6 +414,8 @@ Global
{8DD4F1E8-8AC8-4AB4-AC08-9456AA6EAF10} = {5837821B-CA71-40B6-A9F1-C25D318B4691}
{F7181887-E6B5-4DA9-9598-8E9E806B5A20} = {5837821B-CA71-40B6-A9F1-C25D318B4691}
{922F1EB6-5C8F-45DC-82A8-C651E4554542} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
{115DEB38-679F-467E-8F87-9E669DFD2EB8} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180}
{A855C6B9-1548-4183-926C-75D80CEEF510} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158}

View file

@ -0,0 +1,29 @@
using System;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
using Rebus.Bus;
namespace Elsa.Activities.Rebus
{
[Trigger(Category = "Rebus", Description = "Triggered when a message is received.", Outcomes = new[] { OutcomeNames.Done })]
public class MessageReceived : Activity
{
private readonly IBus _bus;
public MessageReceived(IBus bus)
{
_bus = bus;
}
[ActivityProperty(Hint = "The type of message to receive.")]
public Type MessageType { get; set; } = default!;
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context)
{
var message = context.Input;
return Done(message);
}
}
}

View file

@ -0,0 +1,34 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
using Rebus.Bus;
namespace Elsa.Activities.Rebus
{
[Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })]
public class PublishMessage : Activity
{
private readonly IBus _bus;
public PublishMessage(IBus bus)
{
_bus = bus;
}
[ActivityProperty(Hint = "The message to publish.")]
public object Message { get; set; } = default!;
[ActivityProperty(Hint = "Optional headers to send along with the message.")]
public IDictionary<string, string>? Headers { get; set; }
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
await _bus.Publish(Message, Headers);
return Done();
}
}
}

View file

@ -0,0 +1,37 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
using Rebus.Bus;
namespace Elsa.Activities.Rebus
{
[Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })]
public class SendMessage : Activity
{
private readonly IBus _bus;
public SendMessage(IBus bus)
{
_bus = bus;
}
[ActivityProperty(Hint = "The message to send.")]
public object Message { get; set; } = default!;
[ActivityProperty(Hint = "Optional headers to send along with the message.")]
public IDictionary<string, string>? Headers { get; set; }
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
await _bus.Advanced.Routing.Send("greeting", Message, Headers);
await _bus.Advanced.Topics.Subscribe("greeting");
//await _bus.Send(Message, Headers);
return Done();
}
}
}

View file

@ -0,0 +1,29 @@
using System.Threading.Tasks;
using Elsa.Activities.Rebus.Triggers;
using Elsa.Services;
using Rebus.Extensions;
using Rebus.Handlers;
using Rebus.Messages;
using Rebus.Pipeline;
namespace Elsa.Activities.Rebus.Consumers
{
public class MessageConsumer<T> : IHandleMessages<T>
{
private readonly IWorkflowScheduler _workflowScheduler;
public MessageConsumer(IWorkflowScheduler workflowScheduler)
{
_workflowScheduler = workflowScheduler;
}
public async Task Handle(T message)
{
var correlationId = MessageContext.Current.TransportMessage.Headers.GetValueOrNull(Headers.CorrelationId);
await _workflowScheduler.TriggerWorkflowsAsync<MessageReceivedTrigger>(
x => x.MessageType == typeof(T).Name && (x.CorrelationId == null || x.CorrelationId == correlationId),
message,
correlationId);
}
}
}

View file

@ -0,0 +1,13 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>netstandard2.0</TargetFramework>
<LangVersion>latest</LangVersion>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,2 @@
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>

View file

@ -0,0 +1,26 @@
using Elsa.Activities.Rebus.Consumers;
using Elsa.Activities.Rebus.Triggers;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Activities.Rebus.Extensions
{
public static class ServiceCollectionExtensions
{
public static IServiceCollection AddRebusActivities(this IServiceCollection services) =>
services
.AddTriggerProvider<MessageReceivedTriggerProvider>()
.AddActivity<PublishMessage>()
.AddActivity<SendMessage>()
.AddActivity<MessageReceived>();
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, 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, T5>(this IServiceCollection services) =>
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>>();
}
}

View file

@ -0,0 +1,22 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Triggers;
namespace Elsa.Activities.Rebus.Triggers
{
public class MessageReceivedTrigger : Trigger
{
public string MessageType { get; set; } = default!;
public string? CorrelationId { get; set; }
}
public class MessageReceivedTriggerProvider : TriggerProvider<MessageReceivedTrigger, MessageReceived>
{
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<MessageReceived> context, CancellationToken cancellationToken) =>
new MessageReceivedTrigger
{
MessageType = (await context.Activity.GetPropertyValueAsync(x => x.MessageType, cancellationToken)).Name,
CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId
};
}
}

View file

@ -34,9 +34,6 @@ namespace Elsa.Activities.Timers
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context)
{
if (context.WorkflowExecutionContext.IsFirstPass)
return Done();
ExecuteAt = GetNextOccurrence(CronExpression);
return Suspend();
}

View file

@ -20,17 +20,14 @@ namespace Elsa.Activities.Timers
[ActivityProperty(Hint = "An expression that evaluates to a Duration value.")]
public Duration Timeout { get; set; } = default!;
public Instant ExecuteAt
public Instant? ExecuteAt
{
get => GetState<Instant>();
get => GetState<Instant?>();
set => SetState(value);
}
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context)
{
if (context.WorkflowExecutionContext.IsFirstPass)
return Done();
ExecuteAt = _clock.GetCurrentInstant().Plus(Timeout);
return Suspend();
}

View file

@ -1,7 +1,6 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Triggers;
using NCrontab;
using NodaTime;
namespace Elsa.Activities.Timers.Triggers

View file

@ -1,4 +1,6 @@
using Elsa.Triggers;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Triggers;
using NodaTime;
namespace Elsa.Activities.Timers.Triggers
@ -10,13 +12,27 @@ namespace Elsa.Activities.Timers.Triggers
public class TimerEventTriggerProvider : TriggerProvider<TimerEventTrigger, TimerEvent>
{
public override ITrigger GetTrigger(TriggerProviderContext<TimerEvent> context)
private readonly IClock _clock;
public TimerEventTriggerProvider(IClock clock)
{
var executeAt = context.GetActivity<TimerEvent>().GetState(x => x.ExecuteAt);
_clock = clock;
}
public override async ValueTask<ITrigger> GetTriggerAsync(TriggerProviderContext<TimerEvent> context, CancellationToken cancellationToken)
{
var activity = context.GetActivity<TimerEvent>();
var executeAt = activity.GetState(x => x.ExecuteAt);
if (executeAt == null)
{
var timeout = await activity.GetPropertyValueAsync(x => x.Timeout, cancellationToken);
executeAt = _clock.GetCurrentInstant().Plus(timeout);
}
return new TimerEventTrigger
{
ExecuteAt = executeAt
ExecuteAt = executeAt.Value
};
}
}

View file

@ -1,13 +1,18 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Builders;
using Elsa.Models;
using Elsa.Services.Models;
using Elsa.Triggers;
namespace Elsa.Services
{
public interface IWorkflowRunner
{
Task TriggerWorkflowsAsync<TTrigger>(Func<TTrigger, bool> predicate, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default)
where TTrigger : ITrigger;
ValueTask<WorkflowInstance> RunWorkflowAsync(
WorkflowInstance workflowInstance,
string? activityId = default,
@ -49,7 +54,7 @@ namespace Elsa.Services
string? correlationId = default,
string? contextId = default,
CancellationToken cancellationToken = default);
ValueTask<WorkflowInstance> RunWorkflowAsync(
IWorkflow workflow,
WorkflowInstance workflowInstance,

View file

@ -47,7 +47,7 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton(options.SignalFactory)
.AddSingleton(options.StorageFactory)
.AddPersistence(options.ConfigurePersistence);
options.AddWorkflowsCore();
options.AddMediatR();
options.AddServiceBus();
@ -70,7 +70,7 @@ namespace Microsoft.Extensions.DependencyInjection
.AddTransient<T>()
.AddTransient<IWorkflow>(sp => sp.GetRequiredService<T>());
}
public static IServiceCollection AddWorkflow(this IServiceCollection services, IWorkflow workflow)
{
return services
@ -78,7 +78,14 @@ namespace Microsoft.Extensions.DependencyInjection
.AddTransient(sp => workflow);
}
public static IServiceCollection AddConsumer<TMessage, TConsumer>(this IServiceCollection services) where TConsumer : class, IHandleMessages<TMessage> => services.AddTransient<IHandleMessages<TMessage>, TConsumer>();
public static IServiceCollection AddConsumer<TMessage, TConsumer>(this IServiceCollection services) where TConsumer : class, IHandleMessages<TMessage>
{
return services
.AddTransient<TConsumer>()
.AddTransient<IHandleMessages>(sp => sp.GetRequiredService<TConsumer>())
.AddTransient<IHandleMessages<TMessage>, TConsumer>(sp => sp.GetRequiredService<TConsumer>());
}
private static IServiceCollection AddMediatR(this ElsaOptions options) => options.Services.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity));
private static ElsaOptions AddWorkflowsCore(this ElsaOptions configuration)

View file

@ -6,11 +6,14 @@ using Elsa.ActivityResults;
using Elsa.Builders;
using Elsa.Events;
using Elsa.Exceptions;
using Elsa.Extensions;
using Elsa.Models;
using Elsa.Services.Models;
using Elsa.Triggers;
using MediatR;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Open.Linq.AsyncExtensions;
namespace Elsa.Services
{
@ -23,6 +26,8 @@ namespace Elsa.Services
private readonly IWorkflowRegistry _workflowRegistry;
private readonly IWorkflowFactory _workflowFactory;
private readonly IWorkflowSelector _workflowSelector;
private readonly IWorkflowInstanceManager _workflowInstanceManager;
private readonly Func<IWorkflowBuilder> _workflowBuilderFactory;
private readonly IWorkflowContextManager _workflowContextManager;
private readonly IMediator _mediator;
@ -32,11 +37,12 @@ namespace Elsa.Services
public WorkflowRunner(
IWorkflowRegistry workflowRegistry,
IWorkflowFactory workflowFactory,
IWorkflowSelector workflowSelector,
Func<IWorkflowBuilder> workflowBuilderFactory,
IWorkflowContextManager workflowContextManager,
IMediator mediator,
IServiceProvider serviceProvider,
ILogger<WorkflowRunner> logger)
ILogger<WorkflowRunner> logger, IWorkflowInstanceManager workflowInstanceManager)
{
_workflowRegistry = workflowRegistry;
_workflowFactory = workflowFactory;
@ -45,6 +51,28 @@ namespace Elsa.Services
_mediator = mediator;
_serviceProvider = serviceProvider;
_logger = logger;
_workflowInstanceManager = workflowInstanceManager;
_workflowSelector = workflowSelector;
}
public async Task TriggerWorkflowsAsync<TTrigger>(Func<TTrigger, bool> predicate, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default)
where TTrigger : ITrigger
{
var results = await _workflowSelector.SelectWorkflowsAsync(predicate, cancellationToken).ToList();
foreach (var result in results)
{
if (result.WorkflowInstanceId != null)
{
var workflowInstance = await _workflowInstanceManager.GetByIdAsync(result.WorkflowInstanceId, cancellationToken);
await RunWorkflowAsync(result.WorkflowBlueprint, workflowInstance, result.ActivityId, input, cancellationToken);
}
else
await RunWorkflowAsync(result.WorkflowBlueprint, result.ActivityId, input, correlationId, contextId, cancellationToken);
if (result.Trigger.IsOneOff)
await _workflowSelector.RemoveTriggerAsync(result.Trigger, cancellationToken);
}
}
public async ValueTask<WorkflowInstance> RunWorkflowAsync(
@ -146,7 +174,7 @@ namespace Elsa.Services
await ResumeWorkflowAsync(workflowExecutionContext, activity!, input, cancellationToken);
break;
}
workflowInstance.ContextId = await SaveWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Burst, false, cancellationToken);
await _mediator.Publish(new WorkflowExecuted(workflowExecutionContext), cancellationToken);
@ -173,11 +201,11 @@ namespace Elsa.Services
var context = new LoadWorkflowContext(workflowBlueprint, workflowInstance);
return await _workflowContextManager.LoadContext(context, cancellationToken);
}
private async ValueTask<string?> SaveWorkflowContextAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowContextFidelity fidelity, bool always, CancellationToken cancellationToken)
{
var workflowContext = workflowExecutionContext.WorkflowContext;
if (!always && (workflowContext == null || workflowExecutionContext.WorkflowBlueprint.ContextOptions?.ContextFidelity != fidelity))
return workflowExecutionContext.WorkflowInstance.ContextId;
@ -230,16 +258,16 @@ namespace Elsa.Services
var serviceProvider = scope.ServiceProvider;
var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint;
var workflowInstance = workflowExecutionContext.WorkflowInstance;
while (workflowExecutionContext.HasScheduledActivities)
{
var scheduledActivity = workflowExecutionContext.PopScheduledActivity();
var currentActivityId = scheduledActivity.ActivityId;
var activityBlueprint = workflowBlueprint.GetActivity(currentActivityId)!;
if(workflowBlueprint.ContextOptions?.ContextFidelity == WorkflowContextFidelity.Activity || activityBlueprint.LoadWorkflowContext)
if (workflowBlueprint.ContextOptions?.ContextFidelity == WorkflowContextFidelity.Activity || activityBlueprint.LoadWorkflowContext)
workflowExecutionContext.WorkflowContext = await LoadWorkflowContextAsync(workflowBlueprint, workflowInstance, WorkflowContextFidelity.Activity, activityBlueprint.LoadWorkflowContext, cancellationToken);
var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, serviceProvider, activityBlueprint, scheduledActivity.Input);
var activity = await activityBlueprint.CreateActivityAsync(activityExecutionContext, cancellationToken);
var result = await activityOperation(activityExecutionContext, activity, cancellationToken);

View file

@ -1,8 +1,10 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Messages;
using Elsa.Services;
using Microsoft.Extensions.DependencyInjection;
using Rebus.Handlers;
using Rebus.ServiceProvider;
namespace Elsa.StartupTasks
@ -11,12 +13,19 @@ namespace Elsa.StartupTasks
{
private readonly IServiceProvider _serviceProvider;
public StartServiceBusTask(IServiceProvider serviceProvider) => _serviceProvider = serviceProvider;
public Task ExecuteAsync(CancellationToken cancellationToken = default)
{
_serviceProvider.UseRebus(x => x.Subscribe<RunWorkflow>());
var consumers = _serviceProvider.GetServices<IHandleMessages>();
var messageTypes = consumers.Select(consumer => consumer.GetType().GetInterfaces().First(x => x.GenericTypeArguments.Any()).GenericTypeArguments.First());
_serviceProvider.UseRebus(async bus =>
{
foreach (var messageType in messageTypes)
await bus.Subscribe(messageType);
});
return Task.CompletedTask;
}
}
}

View file

@ -2,7 +2,6 @@
using System.Collections.Generic;
using System.Collections.Immutable;
using System.Linq;
using Elsa.Client.Models;
namespace ElsaDashboard.Application.Models
{

View file

@ -1,6 +1,5 @@
using System;
using System.Threading.Tasks;
using ElsaDashboard.Application.Extensions;
using ElsaDashboard.Application.Models;
using Microsoft.AspNetCore.Components;

View file

@ -1,8 +1,4 @@
using Elsa.Client.Models;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
using NodaTime;
using NodaTime.Serialization.JsonNet;
using ProtoBuf;
namespace ElsaDashboard.Shared.Surrogates

View file

@ -3,7 +3,6 @@ using Elsa.Activities.ControlFlow;
using Elsa.Activities.Timers;
using Elsa.Builders;
using Elsa.Services.Models;
using Microsoft.AspNetCore.Mvc.Filters;
using NodaTime;
namespace Elsa.Samples.ForkJoinTimerAndSignalHttp.Workflows

View file

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

View file

@ -0,0 +1,17 @@
<Project Sdk="Microsoft.NET.Sdk.Worker">
<PropertyGroup>
<TargetFramework>net5.0</TargetFramework>
<UserSecretsId>dotnet-Elsa.Samples.Rebus.AzureServiceBusWorker-34DD18B8-38B1-4F6F-ABE0-BCE3BEEB5FA6</UserSecretsId>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.Hosting" Version="5.0.0" />
<PackageReference Include="Rebus.AzureServiceBus" Version="7.1.6" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\activities\Elsa.Activities.Rebus\Elsa.Activities.Rebus.csproj" />
<ProjectReference Include="..\..\core\Elsa\Elsa.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,9 @@
namespace Elsa.Samples.RebusWorker.Messages
{
public class Greeting
{
public string From { get; set; }
public string To { get; set; }
public string Message { get; set; }
}
}

View file

@ -0,0 +1,50 @@
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
{
public class Program
{
public static void Main(string[] args)
{
CreateHostBuilder(args).Build().Run();
}
public static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureServices((hostContext, services) =>
{
services
.AddElsa(option => option
.UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared"))
.ConfigureServiceBus(ConfigureRebus))
.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

@ -0,0 +1,11 @@
{
"profiles": {
"Elsa.Samples.Rebus.AzureServiceBusWorker": {
"commandName": "Project",
"dotnetRunMessages": "true",
"environmentVariables": {
"DOTNET_ENVIRONMENT": "Development"
}
}
}
}

View file

@ -0,0 +1,21 @@
using Elsa.Activities.Console;
using Elsa.Activities.Rebus;
using Elsa.Builders;
using Elsa.Samples.RebusWorker.Messages;
namespace Elsa.Samples.RebusWorker.Workflows
{
public class ConsumerWorkflow : IWorkflow
{
public void Build(IWorkflowBuilder workflow)
{
workflow
.StartWith<MessageReceived>(messageReceived => messageReceived.Set(x => x.MessageType, typeof(Greeting)))
.WriteLine(context =>
{
var greeting = context.GetInput<Greeting>();
return $"Received a greeting from {greeting.From}, saying \"{greeting.Message}\" to {greeting.To}!";
});
}
}
}

View file

@ -0,0 +1,60 @@
using System;
using Elsa.Activities.Console;
using Elsa.Activities.Rebus;
using Elsa.Activities.Timers;
using Elsa.Builders;
using Elsa.Samples.RebusWorker.Messages;
using NodaTime;
namespace Elsa.Samples.RebusWorker.Workflows
{
public class ProducerWorkflow : IWorkflow
{
private readonly IClock _clock;
private readonly Random _random;
public ProducerWorkflow(IClock clock)
{
_clock = clock;
_random = new Random();
}
public void Build(IWorkflowBuilder workflow)
{
workflow
.TimerEvent(Duration.FromSeconds(5))
.WriteLine("Sending a random greeting to the \"greetings\" queue.")
.Then<PublishMessage>(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting))
//.Then<SendMessage>(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting))
.WriteLine(() => $"Message sent at {_clock.GetCurrentInstant()}");
}
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 index = _random.Next(0, greetings.Length);
return greetings[index];
}
}
}

View file

@ -0,0 +1,9 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
}
}

View file

@ -0,0 +1,9 @@
{
"Logging": {
"LogLevel": {
"Default": "Information",
"Microsoft": "Warning",
"Microsoft.Hosting.Lifetime": "Information"
}
}
}

View file

@ -6,7 +6,6 @@ using AutoFixture;
using Elsa.Activities.Console;
using Elsa.ComponentTests.Helpers;
using Elsa.Models;
using Elsa.Server.Api.Endpoints.WorkflowDefinitions;
using Elsa.Testing.Shared.AutoFixture;
using Elsa.Testing.Shared.Helpers;
using Xunit;