Add workflow dispatch notifications
- Created WorkflowDefinitionDispatching and WorkflowDefinitionDispatched notifications - Created WorkflowInstanceDispatching and WorkflowInstanceDispatched notifications - Updated BackgroundWorkflowDispatcher to emit notifications before and after dispatch - Added integration tests to verify notifications are emitted correctly Co-authored-by: KnibbsyMan <23156317+KnibbsyMan@users.noreply.github.com>
This commit is contained in:
parent
6fc21c5590
commit
f834b040f9
|
|
@ -0,0 +1,10 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Elsa.Workflows.Runtime.Responses;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Notifications;
|
||||
|
||||
/// <summary>
|
||||
/// A notification that is published when a workflow definition has been dispatched.
|
||||
/// </summary>
|
||||
public record WorkflowDefinitionDispatched(DispatchWorkflowDefinitionRequest Request, DispatchWorkflowResponse Response) : INotification;
|
||||
|
|
@ -0,0 +1,9 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Notifications;
|
||||
|
||||
/// <summary>
|
||||
/// A notification that is published when a workflow definition is being dispatched.
|
||||
/// </summary>
|
||||
public record WorkflowDefinitionDispatching(DispatchWorkflowDefinitionRequest Request) : INotification;
|
||||
|
|
@ -0,0 +1,10 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Elsa.Workflows.Runtime.Responses;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Notifications;
|
||||
|
||||
/// <summary>
|
||||
/// A notification that is published when a workflow instance has been dispatched.
|
||||
/// </summary>
|
||||
public record WorkflowInstanceDispatched(DispatchWorkflowInstanceRequest Request, DispatchWorkflowResponse Response) : INotification;
|
||||
|
|
@ -0,0 +1,9 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Notifications;
|
||||
|
||||
/// <summary>
|
||||
/// A notification that is published when a workflow instance is being dispatched.
|
||||
/// </summary>
|
||||
public record WorkflowInstanceDispatching(DispatchWorkflowInstanceRequest Request) : INotification;
|
||||
|
|
@ -3,6 +3,7 @@ using Elsa.Mediator;
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Tenants.Mediator;
|
||||
using Elsa.Workflows.Runtime.Commands;
|
||||
using Elsa.Workflows.Runtime.Notifications;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Elsa.Workflows.Runtime.Responses;
|
||||
|
||||
|
|
@ -11,11 +12,14 @@ namespace Elsa.Workflows.Runtime;
|
|||
/// <summary>
|
||||
/// A simple implementation that queues the specified request for workflow execution on a non-durable background worker.
|
||||
/// </summary>
|
||||
public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IWorkflowDispatcher
|
||||
public class BackgroundWorkflowDispatcher(ICommandSender commandSender, INotificationSender notificationSender, ITenantAccessor tenantAccessor) : IWorkflowDispatcher
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<DispatchWorkflowResponse> DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default)
|
||||
{
|
||||
// Emit dispatching notification
|
||||
await notificationSender.SendAsync(new WorkflowDefinitionDispatching(request), cancellationToken);
|
||||
|
||||
var command = new DispatchWorkflowDefinitionCommand(request.DefinitionVersionId)
|
||||
{
|
||||
Input = request.Input,
|
||||
|
|
@ -27,12 +31,20 @@ public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantA
|
|||
};
|
||||
|
||||
await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken);
|
||||
return DispatchWorkflowResponse.Success();
|
||||
var response = DispatchWorkflowResponse.Success();
|
||||
|
||||
// Emit dispatched notification
|
||||
await notificationSender.SendAsync(new WorkflowDefinitionDispatched(request, response), cancellationToken);
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<DispatchWorkflowResponse> DispatchAsync(DispatchWorkflowInstanceRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default)
|
||||
{
|
||||
// Emit dispatching notification
|
||||
await notificationSender.SendAsync(new WorkflowInstanceDispatching(request), cancellationToken);
|
||||
|
||||
var command = new DispatchWorkflowInstanceCommand(request.InstanceId){
|
||||
BookmarkId = request.BookmarkId,
|
||||
ActivityHandle = request.ActivityHandle,
|
||||
|
|
@ -41,7 +53,12 @@ public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantA
|
|||
CorrelationId = request.CorrelationId};
|
||||
|
||||
await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken);
|
||||
return DispatchWorkflowResponse.Success();
|
||||
var response = DispatchWorkflowResponse.Success();
|
||||
|
||||
// Emit dispatched notification
|
||||
await notificationSender.SendAsync(new WorkflowInstanceDispatched(request, response), cancellationToken);
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
|
|
@ -0,0 +1,16 @@
|
|||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Elsa.Workflows.Runtime.Responses;
|
||||
|
||||
namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications;
|
||||
|
||||
public class Spy
|
||||
{
|
||||
public bool WorkflowDefinitionDispatchingWasCalled { get; set; }
|
||||
public bool WorkflowDefinitionDispatchedWasCalled { get; set; }
|
||||
public bool WorkflowInstanceDispatchingWasCalled { get; set; }
|
||||
public bool WorkflowInstanceDispatchedWasCalled { get; set; }
|
||||
|
||||
public DispatchWorkflowDefinitionRequest? CapturedDefinitionRequest { get; set; }
|
||||
public DispatchWorkflowInstanceRequest? CapturedInstanceRequest { get; set; }
|
||||
public DispatchWorkflowResponse? CapturedResponse { get; set; }
|
||||
}
|
||||
|
|
@ -0,0 +1,46 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Runtime.Notifications;
|
||||
|
||||
namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications;
|
||||
|
||||
public class TestHandler :
|
||||
INotificationHandler<WorkflowDefinitionDispatching>,
|
||||
INotificationHandler<WorkflowDefinitionDispatched>,
|
||||
INotificationHandler<WorkflowInstanceDispatching>,
|
||||
INotificationHandler<WorkflowInstanceDispatched>
|
||||
{
|
||||
private readonly Spy _spy;
|
||||
|
||||
public TestHandler(Spy spy)
|
||||
{
|
||||
_spy = spy;
|
||||
}
|
||||
|
||||
public Task HandleAsync(WorkflowDefinitionDispatching notification, CancellationToken cancellationToken)
|
||||
{
|
||||
_spy.WorkflowDefinitionDispatchingWasCalled = true;
|
||||
_spy.CapturedDefinitionRequest = notification.Request;
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public Task HandleAsync(WorkflowDefinitionDispatched notification, CancellationToken cancellationToken)
|
||||
{
|
||||
_spy.WorkflowDefinitionDispatchedWasCalled = true;
|
||||
_spy.CapturedResponse = notification.Response;
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public Task HandleAsync(WorkflowInstanceDispatching notification, CancellationToken cancellationToken)
|
||||
{
|
||||
_spy.WorkflowInstanceDispatchingWasCalled = true;
|
||||
_spy.CapturedInstanceRequest = notification.Request;
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public Task HandleAsync(WorkflowInstanceDispatched notification, CancellationToken cancellationToken)
|
||||
{
|
||||
_spy.WorkflowInstanceDispatchedWasCalled = true;
|
||||
_spy.CapturedResponse = notification.Response;
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,85 @@
|
|||
using Elsa.Testing.Shared;
|
||||
using Elsa.Workflows.Runtime;
|
||||
using Elsa.Workflows.Runtime.Notifications;
|
||||
using Elsa.Workflows.Runtime.Requests;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Xunit.Abstractions;
|
||||
|
||||
namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications;
|
||||
|
||||
public class Tests
|
||||
{
|
||||
private readonly IWorkflowDispatcher _workflowDispatcher;
|
||||
private readonly Spy _spy;
|
||||
|
||||
public Tests(ITestOutputHelper testOutputHelper)
|
||||
{
|
||||
var services = new TestApplicationBuilder(testOutputHelper)
|
||||
.ConfigureServices(s =>
|
||||
{
|
||||
s.AddSingleton<Spy>();
|
||||
s.AddNotificationHandler<TestHandler, WorkflowDefinitionDispatching>();
|
||||
s.AddNotificationHandler<TestHandler, WorkflowDefinitionDispatched>();
|
||||
s.AddNotificationHandler<TestHandler, WorkflowInstanceDispatching>();
|
||||
s.AddNotificationHandler<TestHandler, WorkflowInstanceDispatched>();
|
||||
})
|
||||
.Build();
|
||||
|
||||
_workflowDispatcher = services.GetRequiredService<IWorkflowDispatcher>();
|
||||
_spy = services.GetRequiredService<Spy>();
|
||||
}
|
||||
|
||||
[Fact(DisplayName = "Dispatching workflow definition should emit notifications")]
|
||||
public async Task DispatchWorkflowDefinition_ShouldEmitNotifications()
|
||||
{
|
||||
// Arrange
|
||||
var definitionVersionId = "test-definition-version-id";
|
||||
var request = new DispatchWorkflowDefinitionRequest(definitionVersionId)
|
||||
{
|
||||
CorrelationId = "test-correlation-id",
|
||||
Input = new Dictionary<string, object> { { "TestKey", "TestValue" } }
|
||||
};
|
||||
|
||||
// Act
|
||||
await _workflowDispatcher.DispatchAsync(request, null);
|
||||
|
||||
// Allow async notification handlers to complete
|
||||
await Task.Delay(100);
|
||||
|
||||
// Assert
|
||||
Assert.True(_spy.WorkflowDefinitionDispatchingWasCalled, "WorkflowDefinitionDispatching notification should be called");
|
||||
Assert.True(_spy.WorkflowDefinitionDispatchedWasCalled, "WorkflowDefinitionDispatched notification should be called");
|
||||
Assert.NotNull(_spy.CapturedDefinitionRequest);
|
||||
Assert.Equal(definitionVersionId, _spy.CapturedDefinitionRequest.DefinitionVersionId);
|
||||
Assert.Equal("test-correlation-id", _spy.CapturedDefinitionRequest.CorrelationId);
|
||||
Assert.NotNull(_spy.CapturedResponse);
|
||||
Assert.True(_spy.CapturedResponse.Succeeded);
|
||||
}
|
||||
|
||||
[Fact(DisplayName = "Dispatching workflow instance should emit notifications")]
|
||||
public async Task DispatchWorkflowInstance_ShouldEmitNotifications()
|
||||
{
|
||||
// Arrange
|
||||
var instanceId = "test-instance-id";
|
||||
var request = new DispatchWorkflowInstanceRequest(instanceId)
|
||||
{
|
||||
CorrelationId = "test-correlation-id",
|
||||
Input = new Dictionary<string, object> { { "TestKey", "TestValue" } }
|
||||
};
|
||||
|
||||
// Act
|
||||
await _workflowDispatcher.DispatchAsync(request, null);
|
||||
|
||||
// Allow async notification handlers to complete
|
||||
await Task.Delay(100);
|
||||
|
||||
// Assert
|
||||
Assert.True(_spy.WorkflowInstanceDispatchingWasCalled, "WorkflowInstanceDispatching notification should be called");
|
||||
Assert.True(_spy.WorkflowInstanceDispatchedWasCalled, "WorkflowInstanceDispatched notification should be called");
|
||||
Assert.NotNull(_spy.CapturedInstanceRequest);
|
||||
Assert.Equal(instanceId, _spy.CapturedInstanceRequest.InstanceId);
|
||||
Assert.Equal("test-correlation-id", _spy.CapturedInstanceRequest.CorrelationId);
|
||||
Assert.NotNull(_spy.CapturedResponse);
|
||||
Assert.True(_spy.CapturedResponse.Succeeded);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,14 @@
|
|||
using Elsa.Workflows.Activities;
|
||||
|
||||
namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications;
|
||||
|
||||
public class SimpleWorkflow : WorkflowBase
|
||||
{
|
||||
public static string DefinitionId = "SimpleWorkflow";
|
||||
|
||||
protected override void Build(IWorkflowBuilder builder)
|
||||
{
|
||||
builder.DefinitionId = DefinitionId;
|
||||
builder.Root = new WriteLine("Hello from SimpleWorkflow");
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue