diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs index d62958a57..a6811a585 100644 --- a/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs @@ -40,6 +40,7 @@ namespace Elsa.Extensions { workflowInstance.Variables = workflowExecutionContext.Variables; workflowInstance.ScheduledActivities = new Stack(workflowExecutionContext.ScheduledActivities.Select(x => new Models.ScheduledActivity(x.Activity.Id, x.Input))); + workflowInstance.Activities = workflowExecutionContext.Activities.Select(x => new ActivityInstance(x.Id, x.Type, x.State, x.Output)).ToList(); workflowInstance.BlockingActivities = new HashSet(workflowExecutionContext.BlockingActivities.Select(x => new BlockingActivity(x.Id, x.Type)), new BlockingActivityEqualityComparer()); workflowInstance.Status = workflowExecutionContext.Status; workflowInstance.CorrelationId = workflowExecutionContext.CorrelationId; diff --git a/src/core/Elsa.Abstractions/Messages/Domain/ActivityExecuted.cs b/src/core/Elsa.Abstractions/Messages/Domain/ActivityExecuted.cs index edd9ece8a..38c04544e 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/ActivityExecuted.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/ActivityExecuted.cs @@ -1,19 +1,11 @@ -using Elsa.Services; -using Elsa.Services.Models; -using MediatR; +using Elsa.Services.Models; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { - public class ActivityExecuted : INotification + public class ActivityExecuted : ActivityNotification { - public ActivityExecuted(WorkflowExecutionContext workflowExecutionContext, ActivityExecutionContext activityExecutionContext) + public ActivityExecuted(ActivityExecutionContext activityExecutionContext) : base(activityExecutionContext) { - WorkflowExecutionContext = workflowExecutionContext; - ActivityExecutionContext = activityExecutionContext; } - - public WorkflowExecutionContext WorkflowExecutionContext { get; } - public ActivityExecutionContext ActivityExecutionContext { get; } - public IActivity Activity => ActivityExecutionContext.Activity; } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Messages/Domain/ActivityExecuting.cs b/src/core/Elsa.Abstractions/Messages/Domain/ActivityExecuting.cs new file mode 100644 index 000000000..4d942e1d7 --- /dev/null +++ b/src/core/Elsa.Abstractions/Messages/Domain/ActivityExecuting.cs @@ -0,0 +1,11 @@ +using Elsa.Services.Models; + +namespace Elsa.Messages.Domain +{ + public class ActivityExecuting : ActivityNotification + { + public ActivityExecuting(ActivityExecutionContext activityExecutionContext) : base(activityExecutionContext) + { + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Messages/Domain/ActivityNotification.cs b/src/core/Elsa.Abstractions/Messages/Domain/ActivityNotification.cs new file mode 100644 index 000000000..ced19f92c --- /dev/null +++ b/src/core/Elsa.Abstractions/Messages/Domain/ActivityNotification.cs @@ -0,0 +1,18 @@ +using Elsa.Services; +using Elsa.Services.Models; +using MediatR; + +namespace Elsa.Messages.Domain +{ + public abstract class ActivityNotification : INotification + { + protected ActivityNotification(ActivityExecutionContext activityExecutionContext) + { + ActivityExecutionContext = activityExecutionContext; + } + + public ActivityExecutionContext ActivityExecutionContext { get; } + public WorkflowExecutionContext WorkflowExecutionContext => ActivityExecutionContext.WorkflowExecutionContext; + public IActivity Activity => ActivityExecutionContext.Activity; + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Messages/Domain/ExecutingActivity.cs b/src/core/Elsa.Abstractions/Messages/Domain/ExecutingActivity.cs deleted file mode 100644 index 385b5bce1..000000000 --- a/src/core/Elsa.Abstractions/Messages/Domain/ExecutingActivity.cs +++ /dev/null @@ -1,14 +0,0 @@ -using Elsa.Services.Models; - -namespace Elsa.Messages -{ - public class ExecutingActivity : WorkflowNotification - { - public ExecutingActivity(WorkflowExecutionContext workflowExecutionContext, ActivityExecutionContext activityExecutionContext) : base(workflowExecutionContext) - { - ActivityExecutionContext = activityExecutionContext; - } - - public ActivityExecutionContext ActivityExecutionContext { get; } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Messages/Domain/ExecutingWorkflow.cs b/src/core/Elsa.Abstractions/Messages/Domain/ExecutingWorkflow.cs index afeaeadd9..5518cb501 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/ExecutingWorkflow.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/ExecutingWorkflow.cs @@ -1,6 +1,6 @@ using Elsa.Services.Models; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { public class ExecutingWorkflow : WorkflowNotification { diff --git a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCancelled.cs b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCancelled.cs index 21a7be0d5..7a5e6d4a3 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCancelled.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCancelled.cs @@ -1,6 +1,6 @@ using Elsa.Services.Models; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { /// /// Published when a workflow transitioned into the Cancelled state. diff --git a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCompleted.cs b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCompleted.cs index 1ea8677fc..4229d78a3 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCompleted.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowCompleted.cs @@ -1,6 +1,6 @@ using Elsa.Services.Models; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { /// /// Published when a workflow transitioned into the Completed state. diff --git a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowDefinitionStoreUpdated.cs b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowDefinitionStoreUpdated.cs index 21fd09225..5f4de9a12 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowDefinitionStoreUpdated.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowDefinitionStoreUpdated.cs @@ -1,6 +1,6 @@ using MediatR; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { /// /// Published when the workflow definition store is updated. diff --git a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowExecuted.cs b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowExecuted.cs index 9d217333c..ba2df1f73 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowExecuted.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowExecuted.cs @@ -1,6 +1,6 @@ using Elsa.Services.Models; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { /// /// Published when a burst of execution finished. diff --git a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowFaulted.cs b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowFaulted.cs index 041f940b5..5f7c5ba08 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowFaulted.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowFaulted.cs @@ -1,6 +1,6 @@ using Elsa.Services.Models; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { /// /// Published when a workflow transitioned into the Faulted state. diff --git a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowNotification.cs b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowNotification.cs index 9c4ec8bce..805263dcf 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowNotification.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowNotification.cs @@ -1,7 +1,7 @@ using Elsa.Services.Models; using MediatR; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { /// /// Common base for workflow-related events. diff --git a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowSuspended.cs b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowSuspended.cs index a0badc625..a1dc0baf0 100644 --- a/src/core/Elsa.Abstractions/Messages/Domain/WorkflowSuspended.cs +++ b/src/core/Elsa.Abstractions/Messages/Domain/WorkflowSuspended.cs @@ -1,6 +1,6 @@ using Elsa.Services.Models; -namespace Elsa.Messages +namespace Elsa.Messages.Domain { /// /// Published when a workflow transitioned into the Suspended state. diff --git a/src/core/Elsa.Abstractions/Models/ActivityInstance.cs b/src/core/Elsa.Abstractions/Models/ActivityInstance.cs index 580062300..85edc96c6 100644 --- a/src/core/Elsa.Abstractions/Models/ActivityInstance.cs +++ b/src/core/Elsa.Abstractions/Models/ActivityInstance.cs @@ -2,6 +2,18 @@ namespace Elsa.Models { public class ActivityInstance { + public ActivityInstance() + { + } + + public ActivityInstance(string id, string type, Variables state, Variable? output) + { + Id = id; + Type = type; + State = state; + Output = output; + } + public string? Id { get; set; } public string? Type { get; set; } public Variables? State { get; set; } diff --git a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs index 1faa90701..1eb972d01 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs @@ -9,6 +9,7 @@ namespace Elsa.Models public WorkflowInstance() { Variables = new Variables(); + Activities = new List(); BlockingActivities = new HashSet(new BlockingActivityEqualityComparer()); ExecutionLog = new List(); ScheduledActivities = new Stack(); @@ -22,6 +23,7 @@ namespace Elsa.Models public Instant CreatedAt { get; set; } public Variables Variables { get; set; } public Variable? Output { get; set; } + public ICollection Activities { get; set; } public HashSet BlockingActivities { get; set; } public ICollection ExecutionLog { get; set; } public WorkflowFault? Fault { get; set; } diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs b/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs index d90560ea6..4a85098bd 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs @@ -8,7 +8,7 @@ namespace Elsa.Services { Task RunWorkflowAsync(Workflow workflow, string? activityId = default, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default); - Task RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default); + Task RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default); /// /// Run a registered workflow by its ID. diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionScope.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionScope.cs deleted file mode 100644 index 940f2cf01..000000000 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionScope.cs +++ /dev/null @@ -1,18 +0,0 @@ -using Elsa.Models; - -namespace Elsa.Services.Models -{ - public class WorkflowExecutionScope - { - public WorkflowExecutionScope(Variables? variables = default) - { - Variables = variables ?? new Variables(); - } - - public Variables Variables { get; } - - public void SetVariable(string variableName, object value) => Variables.SetVariable(variableName, value); - public T GetVariable(string name) => Variables.GetVariable(name); - public object GetVariable(string name) => Variables.GetVariable(name); - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs b/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs index 37bdf75e2..889017722 100644 --- a/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs +++ b/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs @@ -6,6 +6,7 @@ using Elsa.Attributes; using Elsa.Design; using Elsa.Extensions; using Elsa.Messages; +using Elsa.Messages.Domain; using Elsa.Results; using Elsa.Services; using Elsa.Services.Models; diff --git a/src/core/Elsa.Core/Messages/Domain/Handlers/PersistenceWorkflowEventHandler.cs b/src/core/Elsa.Core/Messages/Domain/Handlers/PersistenceWorkflowEventHandler.cs index 316fce992..623c97fe8 100644 --- a/src/core/Elsa.Core/Messages/Domain/Handlers/PersistenceWorkflowEventHandler.cs +++ b/src/core/Elsa.Core/Messages/Domain/Handlers/PersistenceWorkflowEventHandler.cs @@ -1,57 +1,69 @@ using System; using System.Threading; using System.Threading.Tasks; +using Elsa.Extensions; using Elsa.Models; using Elsa.Persistence; +using Elsa.Services.Models; using MediatR; +using Microsoft.Extensions.Logging; -namespace Elsa.Messages.Handlers +namespace Elsa.Messages.Domain.Handlers { - public class PersistenceWorkflowEventHandler : - INotificationHandler, + public class PersistenceWorkflowEventHandler : + INotificationHandler, INotificationHandler, INotificationHandler, INotificationHandler { private readonly IWorkflowInstanceStore workflowInstanceStore; + private readonly ILogger logger; - public PersistenceWorkflowEventHandler(IWorkflowInstanceStore workflowInstanceStore) + public PersistenceWorkflowEventHandler(IWorkflowInstanceStore workflowInstanceStore, ILogger logger) { this.workflowInstanceStore = workflowInstanceStore; + this.logger = logger; } - + public async Task Handle(WorkflowSuspended notification, CancellationToken cancellationToken) { if (notification.WorkflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.Suspended) - await SaveWorkflowAsync(cancellationToken); + await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken); } - + public async Task Handle(WorkflowExecuted notification, CancellationToken cancellationToken) { if (notification.WorkflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.WorkflowExecuted) - await SaveWorkflowAsync(cancellationToken); + await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken); } - + public async Task Handle(ActivityExecuted notification, CancellationToken cancellationToken) { if (notification.WorkflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.ActivityExecuted) - await SaveWorkflowAsync(cancellationToken); - } - - public async Task Handle(WorkflowCompleted notification, CancellationToken cancellationToken) - { - var workflow = notification.WorkflowExecutionContext; - var blueprint = workflow; - - if (blueprint.DeleteCompletedInstances || blueprint.PersistenceBehavior == WorkflowPersistenceBehavior.Suspended) - await workflowInstanceStore.DeleteAsync(workflow.InstanceId, cancellationToken); + await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken); } - private async Task SaveWorkflowAsync(CancellationToken cancellationToken) + public async Task Handle(WorkflowCompleted notification, CancellationToken cancellationToken) { - //var workflowInstance = workflow.ToInstance(); - //await workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken); - throw new NotImplementedException(); + var workflowExecutionContext = notification.WorkflowExecutionContext; + + if (workflowExecutionContext.DeleteCompletedInstances || workflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.Suspended) + { + logger.LogDebug("Deleting completed workflow instance {WorkflowInstanceId}", workflowExecutionContext.InstanceId); + await workflowInstanceStore.DeleteAsync(workflowExecutionContext.InstanceId, cancellationToken); + } + } + + private async Task SaveWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken) + { + var workflowInstance = await workflowInstanceStore.GetByIdAsync(workflowExecutionContext.InstanceId, cancellationToken); + + if (workflowInstance == null) + workflowInstance = workflowExecutionContext.CreateWorkflowInstance(); + else + workflowInstance = workflowExecutionContext.UpdateWorkflowInstance(workflowInstance); + + await workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Persistence/PublishingWorkflowDefinitionStore.cs b/src/core/Elsa.Core/Persistence/PublishingWorkflowDefinitionStore.cs index a15a7c114..124ccaf13 100644 --- a/src/core/Elsa.Core/Persistence/PublishingWorkflowDefinitionStore.cs +++ b/src/core/Elsa.Core/Persistence/PublishingWorkflowDefinitionStore.cs @@ -2,6 +2,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Messages; +using Elsa.Messages.Domain; using Elsa.Models; using MediatR; diff --git a/src/core/Elsa.Core/Runtime/ServiceProviderExtensions.cs b/src/core/Elsa.Core/Runtime/ServiceProviderExtensions.cs deleted file mode 100644 index e57865696..000000000 --- a/src/core/Elsa.Core/Runtime/ServiceProviderExtensions.cs +++ /dev/null @@ -1,27 +0,0 @@ -using System; -using System.Threading; -using System.Threading.Tasks; -using Elsa.Messages.Distributed; -using Microsoft.Extensions.DependencyInjection; -using Rebus.Bus; -using Rebus.ServiceProvider; - -namespace Elsa.Runtime -{ - public static class ServiceProviderExtensions - { - /// - /// Starts the service bus and registers workflow message consumers. - /// - public static async Task StartElsaAsync(this IServiceProvider serviceProvider, CancellationToken cancellationToken = default) - { - var bus = serviceProvider.GetRequiredService(); - - serviceProvider.UseRebus(); - - await bus.Subscribe(); - - return serviceProvider; - } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowHost.cs b/src/core/Elsa.Core/Services/WorkflowHost.cs index 21d9c3068..99d0c9080 100644 --- a/src/core/Elsa.Core/Services/WorkflowHost.cs +++ b/src/core/Elsa.Core/Services/WorkflowHost.cs @@ -7,11 +7,13 @@ using Elsa.Comparers; using Elsa.Expressions; using Elsa.Extensions; using Elsa.Messages; +using Elsa.Messages.Domain; using Elsa.Models; using Elsa.Persistence; using Elsa.Results; using Elsa.Services.Models; using MediatR; +using Microsoft.Extensions.Logging; using ScheduledActivity = Elsa.Services.Models.ScheduledActivity; namespace Elsa.Services @@ -30,6 +32,7 @@ namespace Elsa.Services private readonly IIdGenerator idGenerator; private readonly IMediator mediator; private readonly IServiceProvider serviceProvider; + private readonly ILogger logger; public WorkflowHost( IWorkflowRegistry workflowRegistry, @@ -38,7 +41,8 @@ namespace Elsa.Services IExpressionEvaluator expressionEvaluator, IIdGenerator idGenerator, IMediator mediator, - IServiceProvider serviceProvider) + IServiceProvider serviceProvider, + ILogger logger) { this.workflowRegistry = workflowRegistry; this.workflowInstanceStore = workflowInstanceStore; @@ -47,11 +51,19 @@ namespace Elsa.Services this.idGenerator = idGenerator; this.mediator = mediator; this.serviceProvider = serviceProvider; + this.logger = logger; } - public async Task RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default) + public async Task RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default) { var workflowInstance = await workflowInstanceStore.GetByIdAsync(workflowInstanceId, cancellationToken); + + if (workflowInstance == null) + { + logger.LogDebug("Workflow instance {WorkflowInstanceId} does not exist.", workflowInstanceId); + return null; + } + var workflow = await workflowRegistry.GetWorkflowAsync(workflowInstance.DefinitionId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken); return await RunAsync(workflow, workflowInstance, activityId, input, cancellationToken); } @@ -127,13 +139,36 @@ namespace Elsa.Services var result = await activityOperation(activityExecutionContext, currentActivity, cancellationToken); await result.ExecuteAsync(activityExecutionContext, cancellationToken); - await mediator.Publish(new ActivityExecuted(workflowExecutionContext, activityExecutionContext), cancellationToken); + await mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken); activityOperation = Execute; } if (workflowExecutionContext.Status == WorkflowStatus.Running) workflowExecutionContext.Complete(); + + await mediator.Publish(new WorkflowExecuted(workflowExecutionContext), cancellationToken); + + var statusEvent = default(object); + + switch (workflowExecutionContext.Status) + { + case WorkflowStatus.Cancelled: + statusEvent = new WorkflowCancelled(workflowExecutionContext); + break; + case WorkflowStatus.Completed: + statusEvent = new WorkflowCompleted(workflowExecutionContext); + break; + case WorkflowStatus.Faulted: + statusEvent = new WorkflowFaulted(workflowExecutionContext); + break; + case WorkflowStatus.Suspended: + statusEvent = new WorkflowSuspended(workflowExecutionContext); + break; + } + + if (statusEvent != null) + await mediator.Publish(statusEvent, cancellationToken); } private ScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel, IDictionary activityLookup) @@ -145,11 +180,23 @@ namespace Elsa.Services private WorkflowExecutionContext CreateWorkflowExecutionContext(Workflow workflow, WorkflowInstance workflowInstance) { var activityLookup = workflow.Activities.ToDictionary(x => x.Id); + var activityInstanceLookup = workflowInstance.Activities.ToDictionary(x => x.Id); var scheduledActivities = new Stack(workflowInstance.ScheduledActivities.Reverse().Select(x => CreateScheduledActivity(x, activityLookup))); var blockingActivities = new HashSet(workflowInstance.BlockingActivities.Select(x => activityLookup[x.ActivityId])); var variables = workflowInstance.Variables; var status = workflowInstance.Status; var persistenceBehavior = workflow.PersistenceBehavior; + + foreach (var activity in workflow.Activities) + { + if (!activityInstanceLookup.ContainsKey(activity.Id)) + continue; + + var activityInstance = activityInstanceLookup[activity.Id]; + activity.State = activityInstance.State; + activity.Output = activityInstance.Output; + } + return CreateWorkflowExecutionContext( workflowInstance.Id, workflow.DefinitionId, diff --git a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs index 75a20cdac..2f7df8477 100644 --- a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs +++ b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs @@ -1,14 +1,15 @@ -using System; using System.Threading; using System.Threading.Tasks; +using Elsa.Messages.Distributed; using Elsa.Runtime; +using Rebus.Bus; namespace Elsa.StartupTasks { public class StartServiceBusTask : IStartupTask { - private readonly IServiceProvider serviceProvider; - public StartServiceBusTask(IServiceProvider serviceProvider) => this.serviceProvider = serviceProvider; - public Task ExecuteAsync(CancellationToken cancellationToken = default) => serviceProvider.StartElsaAsync(); + private readonly IBus serviceBus; + public StartServiceBusTask(IBus serviceBus) => this.serviceBus = serviceBus; + public Task ExecuteAsync(CancellationToken cancellationToken = default) => serviceBus.Subscribe(); } } \ No newline at end of file diff --git a/src/samples/Elsa.Samples.Timers/Program.cs b/src/samples/Elsa.Samples.Timers/Program.cs index 71fecf30c..c488c927d 100644 --- a/src/samples/Elsa.Samples.Timers/Program.cs +++ b/src/samples/Elsa.Samples.Timers/Program.cs @@ -2,6 +2,7 @@ using Elsa.Runtime; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; using NodaTime; namespace Elsa.Samples.Timers @@ -10,13 +11,12 @@ namespace Elsa.Samples.Timers { private static async Task Main() { - var host = CreateHostBuilder().UseConsoleLifetime().Build(); - await host.Services.StartElsaAsync(); - await host.RunAsync(); + await CreateHostBuilder().UseConsoleLifetime().Build().RunAsync(); } public static IHostBuilder CreateHostBuilder() => Host.CreateDefaultBuilder() + .ConfigureLogging(logging => logging.AddConsole().SetMinimumLevel(LogLevel.Debug)) .ConfigureServices((hostContext, services) => { services diff --git a/test/unit/Elsa.Core.UnitTests/WorkflowRunnerTests.cs b/test/unit/Elsa.Core.UnitTests/WorkflowHostTests.cs similarity index 96% rename from test/unit/Elsa.Core.UnitTests/WorkflowRunnerTests.cs rename to test/unit/Elsa.Core.UnitTests/WorkflowHostTests.cs index 1352a92c4..572410e96 100644 --- a/test/unit/Elsa.Core.UnitTests/WorkflowRunnerTests.cs +++ b/test/unit/Elsa.Core.UnitTests/WorkflowHostTests.cs @@ -19,11 +19,11 @@ using Xunit; namespace Elsa.Core.UnitTests { - public class WorkflowRunnerTests + public class WorkflowHostTests { private readonly WorkflowHost workflowHost; - public WorkflowRunnerTests() + public WorkflowHostTests() { var fixture = new Fixture().Customize(new NodaTimeCustomization()); var workflowActivatorMock = new Mock(); @@ -44,7 +44,8 @@ namespace Elsa.Core.UnitTests workflowExpressionEvaluatorMock.Object, idGeneratorMock.Object, mediatorMock.Object, - serviceProvider); + serviceProvider, + logger); } [Fact(DisplayName = "Can run simple workflow to completed state.")]