Rename ServiceBus integration tests project

This commit is contained in:
Sipke Schoorstra 2023-10-20 13:38:40 +02:00
parent d6481133fd
commit bd1eb2edfa
13 changed files with 405 additions and 414 deletions

View file

@ -263,7 +263,7 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Alterations.MassTransi
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Activities.UnitTests", "test\unit\Elsa.Activities.UnitTests\Elsa.Activities.UnitTests.csproj", "{E6562B0F-AF64-472A-B009-BF6D40DDE99B}"
EndProject
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.ServiceBusIntegrationTests", "test\integration\Elsa.ServiceBusIntegrationTests\Elsa.ServiceBusIntegrationTests.csproj", "{9504E3F3-F77F-437A-8644-D23D7F5FCF8B}"
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.ServiceBus.IntegrationTests", "test\integration\Elsa.ServiceBus.IntegrationTests\Elsa.ServiceBus.IntegrationTests.csproj", "{9504E3F3-F77F-437A-8644-D23D7F5FCF8B}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "ElsaStudioWebAssembly", "src\bundles\ElsaStudioWebAssembly\ElsaStudioWebAssembly.csproj", "{E12A1BDF-5D65-493B-835D-AFBC480A5650}"
EndProject

View file

@ -1,6 +1,6 @@
using Elsa.ServiceBusIntegrationTests.Helpers;
using Elsa.ServiceBus.IntegrationTests.Helpers;
namespace Elsa.ServiceBusIntegrationTests.Contracts
namespace Elsa.ServiceBus.IntegrationTests.Contracts
{
public interface IServiceBusProcessorManager
{

View file

@ -1,4 +1,4 @@
namespace Elsa.ServiceBusIntegrationTests.Contracts
namespace Elsa.ServiceBus.IntegrationTests.Contracts
{
public interface ITestResetEventManager
{

View file

@ -1,23 +1,23 @@
using Azure.Messaging.ServiceBus;
using Elsa.ServiceBusIntegrationTests.Contracts;
using Elsa.ServiceBus.IntegrationTests.Contracts;
using NSubstitute;
using Xunit.Abstractions;
namespace Elsa.ServiceBusIntegrationTests.Helpers
namespace Elsa.ServiceBus.IntegrationTests.Helpers
{
public class ServiceBusProcessorManager : IServiceBusProcessorManager
{
private readonly ServiceBusClient _serviceBusClient;
private readonly ITestOutputHelper _testOutputHelper;
private readonly IDictionary<(string topicName, string subscriptionName), ServiceBusProcessorTest> _serviceBusProcessors
= new Dictionary<(string topicName, string subscriptionName), ServiceBusProcessorTest>();
public ServiceBusProcessorManager(ServiceBusClient serviceBusClient, ITestOutputHelper testOutputHelper)
private readonly IDictionary<(string topicName, string subscriptionName), ServiceBusProcessorTest> _serviceBusProcessors = new Dictionary<(string topicName, string subscriptionName), ServiceBusProcessorTest>();
public ServiceBusProcessorManager(ServiceBusClient serviceBusClient, ITestOutputHelper testOutputHelper)
{
_serviceBusClient = serviceBusClient;
_testOutputHelper = testOutputHelper;
}
public ServiceBusProcessorTest Init(string topic ,string subscription)
public ServiceBusProcessorTest Init(string topic, string subscription)
{
_serviceBusProcessors.TryGetValue((topic, subscription), out var processor);
if (processor != null)
@ -28,16 +28,17 @@ namespace Elsa.ServiceBusIntegrationTests.Helpers
.CreateProcessor(Arg.Is<string>(s => s == topic),
Arg.Any<string>(),
Arg.Any<ServiceBusProcessorOptions>()
)
.Returns((c)=> {
)
.Returns((c) =>
{
_testOutputHelper.WriteLine("ServiceBusClient");
return processor;
});
_serviceBusProcessors.Add((topic,subscription), processor);
_serviceBusProcessors.Add((topic, subscription), processor);
return processor;
}
public ServiceBusProcessorTest Get(string topic, string subscription)
{
_serviceBusProcessors.TryGetValue((topic, subscription), out var processor);

View file

@ -2,7 +2,7 @@
using System.Text.Json;
using Xunit.Abstractions;
namespace Elsa.ServiceBusIntegrationTests.Helpers
namespace Elsa.ServiceBus.IntegrationTests.Helpers
{
public class ServiceBusProcessorTest : ServiceBusProcessor
{
@ -17,8 +17,14 @@ namespace Elsa.ServiceBusIntegrationTests.Helpers
var args = CreateMessageArgs(payload, correlationId, attempt);
await base.OnProcessMessageAsync(args);
}
public ProcessMessageEventArgs CreateMessageArgs<T>(T payload, string correlationId, int deliveryCount = 1)
public override Task StartProcessingAsync(CancellationToken cancellationToken = default)
{
_testOutputHelper.WriteLine("Receiving Service Bus Message");
return Task.CompletedTask;
}
private ProcessMessageEventArgs CreateMessageArgs<T>(T payload, string correlationId, int deliveryCount = 1)
{
var payloadJson = JsonSerializer.Serialize(payload);
var props = new Dictionary<string, object>() { };
@ -28,17 +34,13 @@ namespace Elsa.ServiceBusIntegrationTests.Helpers
deliveryCount: deliveryCount,
correlationId: correlationId,
properties: props
);
);
var args = new ProcessMessageEventArgs(message, null, new CancellationToken());
return args;
}
public override async Task StartProcessingAsync(CancellationToken cancellationToken = default)
{
_testOutputHelper.WriteLine("Receiving Service Bus Message");
}
}
}

View file

@ -1,6 +1,6 @@
using Elsa.ServiceBusIntegrationTests.Contracts;
using Elsa.ServiceBus.IntegrationTests.Contracts;
namespace Elsa.ServiceBusIntegrationTests.Helpers
namespace Elsa.ServiceBus.IntegrationTests.Helpers
{
public class TestResetEventManager : ITestResetEventManager
{

View file

@ -0,0 +1,372 @@
using Azure.Messaging.ServiceBus;
using Elsa.AzureServiceBus.Contracts;
using Elsa.AzureServiceBus.Services;
using Elsa.Testing.Shared;
using Elsa.Workflows.Core;
using Microsoft.Extensions.DependencyInjection;
using Xunit.Abstractions;
using Elsa.Mediator.HostedServices;
using Elsa.Mediator.Options;
using Elsa.ServiceBus.IntegrationTests.Contracts;
using Elsa.ServiceBus.IntegrationTests.Helpers;
using Elsa.ServiceBus.IntegrationTests.Scenarios.Workflows;
using Microsoft.Extensions.Options;
using Elsa.Workflows.Runtime.Options;
using Elsa.Workflows.Runtime.Contracts;
using NSubstitute;
namespace Elsa.ServiceBus.IntegrationTests.Scenarios.ServiceBus;
public class ServiceBusTest : IDisposable
{
private readonly IServiceProvider _services;
private readonly CapturingTextWriter _capturingTextWriter = new();
private readonly ServiceBusClient _serviceBusClient = Substitute.For<ServiceBusClient>();
private readonly IWorkerManager _worker;
private readonly BackgroundCommandSenderHostedService _hosted;
private readonly BackgroundEventPublisherHostedService _backgroundEventService;
private readonly ITestResetEventManager _resetEventManager = new TestResetEventManager();
private readonly ITestOutputHelper _testOutputHelper;
private readonly IServiceBusProcessorManager _sbProcessorManager;
public ServiceBusTest(ITestOutputHelper testOutputHelper)
{
_services = new TestApplicationBuilder(testOutputHelper)
.WithCapturingTextWriter(_capturingTextWriter)
.AddWorkflow<ReceiveOneMessageWorkflow>()
.AddWorkflow<ReceiveMessageWorkflow>()
.AddWorkflow<SendOneMessageWorkflow>()
.AddWorkflow<SendOneMessageWithCorrelationIdWorkflow>()
.ConfigureServices(services =>
{
services
.AddSingleton(_serviceBusClient)
.AddSingleton<IServiceBusProcessorManager, ServiceBusProcessorManager>()
.AddSingleton<IWorkerManager, WorkerManager>()
.AddSingleton(_resetEventManager)
.AddSingleton(sp =>
{
var options = sp.GetRequiredService<IOptions<MediatorOptions>>().Value;
return ActivatorUtilities.CreateInstance<BackgroundCommandSenderHostedService>(sp, options.CommandWorkerCount);
})
.AddSingleton(sp =>
{
var options = sp.GetRequiredService<IOptions<MediatorOptions>>().Value;
return ActivatorUtilities.CreateInstance<BackgroundEventPublisherHostedService>(sp, options.NotificationWorkerCount);
})
;
})
.Build();
_worker = _services.GetRequiredService<IWorkerManager>();
_hosted = _services.GetRequiredService<BackgroundCommandSenderHostedService>();
_backgroundEventService = _services.GetRequiredService<BackgroundEventPublisherHostedService>();
_sbProcessorManager = _services.GetRequiredService<IServiceBusProcessorManager>();
_testOutputHelper = testOutputHelper;
}
[Fact(DisplayName = "2 Receive - Sending 1 message - Should Block")]
public async Task Receive_1_Message_Should_Block_If_One_Receive()
{
_sbProcessorManager.Init("topicName", "subscriptionName");
_sbProcessorManager.Init("topicName1", "subscription1");
//Init waitEvent :
_resetEventManager.Init("receive1");
_resetEventManager.Init("receive2");
//Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
//Start Workflow
const string workflowDefinitionId = nameof(ReceiveMessageWorkflow);
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Start Worker to send Message on topicName/subscriptionName
await _worker.StartWorkerAsync("topicName", "subscriptionName");
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, null!);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
//Wait for receiving second message
var wait2 = _resetEventManager.Get("receive2").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait2 : {wait2}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We don't send 2 messages so Workflow must be
* Running
* Suspended
*
* with a timeout on Wait for 2nd Message. (False ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Running, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, lastWorkflowState.SubStatus);
Assert.True(wait1);
Assert.False(wait2);
}
private async Task InitRegistryAndBackGroundServiceWorkerAsync()
{
// Init Registries to use StartWorkflow
await _services.PopulateRegistriesAsync();
// Start background services for CommandHandler
await _hosted.StartAsync(CancellationToken.None);
await _backgroundEventService.StartAsync(CancellationToken.None);
}
[Fact(DisplayName = "1 Receive - Sending 1 message - Should Finished")]
public async Task Receive_1_Message_Should_Finish_With_One_Receive()
{
_sbProcessorManager.Init("topicName", "subscriptionName");
// Init waitEvent:
_resetEventManager.Init("receive1");
// Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
// Start Workflow
const string workflowDefinitionId = nameof(ReceiveOneMessageWorkflow);
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Start Worker to send Message on topicName/subscriptionName
await _worker.StartWorkerAsync("topicName", "subscriptionName");
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, null!);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We sent 1 message so Workflow must be
* Finished
* Finished
*
* with no timeout for the first message. (True ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Finished, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Finished, lastWorkflowState.SubStatus);
Assert.True(wait1);
}
[Fact(DisplayName = "1 Send - 1 Receive - Should Finished - without race condition")]
public async Task Send_1_Message_And_Receive_Response()
{
_sbProcessorManager.Init("topicName", "subscriptionName");
// Init waitEvent:
_resetEventManager.Init("receive1");
// Init Background
await InitRegistryAndBackGroundServiceWorkerAsync();
var senderMock = Substitute.For<ServiceBusSender>();
_serviceBusClient
.CreateSender(Arg.Any<string>())
.Returns(senderMock);
senderMock
.SendMessageAsync(Arg.Any<ServiceBusMessage>(), Arg.Any<CancellationToken>())
.ReturnsForAnyArgs(async (callback) =>
{
var sb = callback.Arg<ServiceBusMessage>();
var c = callback.Arg<CancellationToken>();
_testOutputHelper.WriteLine("Sending Message from activity");
await _worker.StartWorkerAsync("topicName", "subscriptionName", c);
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, null!);
});
// Start Workflow
const string workflowDefinitionId = nameof(SendOneMessageWorkflow);
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We sent 1 message so Workflow must be
* Finished
* Finished
*
* with no timeout for the first message. (True ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Finished, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Finished, lastWorkflowState.SubStatus);
Assert.True(wait1);
}
[Fact(DisplayName = "1 Send - 1 Receive - with correlationId - Should Finished - without race condition")]
public async Task Send_1_Message_And_Correlate_Receive_Response()
{
var correlationId = "EEE3D9CC-2279-4CE5-8F4F-FC6C65BF8814";
_sbProcessorManager.Init("topicName", "subscriptionName");
//Init waitEvent :
_resetEventManager.Init("receive1");
//Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
var senderMock = Substitute.For<ServiceBusSender>();
_serviceBusClient
.CreateSender(Arg.Any<string>())
.Returns(senderMock);
senderMock
.SendMessageAsync(Arg.Any<ServiceBusMessage>(), Arg.Any<CancellationToken>())
.ReturnsForAnyArgs(async (callback) =>
{
var sb = callback.Arg<ServiceBusMessage>();
var c = callback.Arg<CancellationToken>();
_testOutputHelper.WriteLine("Sending Message from activity");
await _worker.StartWorkerAsync("topicName", "subscriptionName");
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, correlationId);
});
//Start Workflow
var workflowDefinitionId = nameof(SendOneMessageWithCorrelationIdWorkflow);
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
// Wait for receiving first message.
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We sent 1 message so Workflow must be
* Finished
* Finished
*
* with no timeout for the first message. (True ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Finished, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Finished, lastWorkflowState.SubStatus);
Assert.True(wait1);
}
[Fact(DisplayName = "1 Receive - Listening from other topic - Should not trigger")]
public async Task Receive_1_Message_Should_Not_Trigger()
{
_sbProcessorManager.Init("topicName1", "subscriptionName1");
// Init waitEvent:
_resetEventManager.Init("receive1");
_resetEventManager.Init("receive2");
// Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
// Start Workflow
var workflowDefinitionId = nameof(ReceiveMessageWorkflow);
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Start Worker to send Message on topicName/suscriptionName
await _worker.StartWorkerAsync("topicName1", "subscriptionName1");
await _sbProcessorManager
.Get("topicName1", "subscriptionName1")
.SendMessage<dynamic>(new { hello = "world" }, null!);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We don't send 2 messages so Workflow must be
* Running
* Suspended
*
* with a timeout on Wait for 2nd Message. (False ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Running, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, lastWorkflowState.SubStatus);
Assert.False(wait1);
}
void IDisposable.Dispose()
{
_hosted.StopAsync(new CancellationToken()).GetAwaiter().GetResult();
}
}

View file

@ -1,10 +1,10 @@
using Elsa.AzureServiceBus.Activities;
using Elsa.ServiceBus.IntegrationTests.Contracts;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Activities;
using Elsa.Workflows.Core.Contracts;
using Elsa.ServiceBusIntegrationTests.Contracts;
namespace Elsa.ServiceBusIntegrationTests.Scenarios.Workflows
namespace Elsa.ServiceBus.IntegrationTests.Scenarios.Workflows
{
public class ReceiveMessageWorkflow : WorkflowBase
{

View file

@ -1,10 +1,10 @@
using Elsa.AzureServiceBus.Activities;
using Elsa.ServiceBus.IntegrationTests.Contracts;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Activities;
using Elsa.Workflows.Core.Contracts;
using Elsa.ServiceBusIntegrationTests.Contracts;
namespace Elsa.ServiceBusIntegrationTests.Scenarios.Workflows
namespace Elsa.ServiceBus.IntegrationTests.Scenarios.Workflows
{
public class ReceiveOneMessageWorkflow : WorkflowBase
{

View file

@ -1,10 +1,10 @@
using Elsa.AzureServiceBus.Activities;
using Elsa.ServiceBus.IntegrationTests.Contracts;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Activities;
using Elsa.Workflows.Core.Contracts;
using Elsa.ServiceBusIntegrationTests.Contracts;
namespace Elsa.ServiceBusIntegrationTests.Scenarios.Workflows
namespace Elsa.ServiceBus.IntegrationTests.Scenarios.Workflows
{
public class SendOneMessageWorkflow : WorkflowBase
{

View file

@ -1,384 +0,0 @@
using Azure.Messaging.ServiceBus;
using Elsa.AzureServiceBus.Contracts;
using Elsa.AzureServiceBus.Models;
using Elsa.AzureServiceBus.Services;
using Elsa.Testing.Shared;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Contracts;
using Microsoft.Extensions.DependencyInjection;
using System.Collections.ObjectModel;
using Xunit.Abstractions;
using Elsa.Mediator.HostedServices;
using Elsa.Mediator.Options;
using Microsoft.Extensions.Options;
using Elsa.Workflows.Runtime.Options;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.ServiceBusIntegrationTests.Contracts;
using Elsa.ServiceBusIntegrationTests.Scenarios.Workflows;
using Elsa.ServiceBusIntegrationTests.Helpers;
using NSubstitute;
namespace Elsa.ServiceBusIntegrationTests.Scenarios.ServiceBus
{
public class ServiceBusTest : IDisposable
{
private readonly IServiceProvider _services;
private readonly CapturingTextWriter _capturingTextWriter = new();
private readonly ServiceBusClient _serviceBusClient = Substitute.For<ServiceBusClient>();
private IWorkerManager _worker;
private readonly BackgroundCommandSenderHostedService _hosted;
private readonly BackgroundEventPublisherHostedService _backgroundEventService;
private readonly ITestResetEventManager _resetEventManager = new TestResetEventManager();
private readonly ITestOutputHelper _testOutputHelper;
private readonly IServiceBusProcessorManager _sbProcessorManager;
public ServiceBusTest(ITestOutputHelper testOutputHelper)
{
_services = new TestApplicationBuilder(testOutputHelper)
.WithCapturingTextWriter(_capturingTextWriter)
.AddWorkflow<ReceiveOneMessageWorkflow>()
.AddWorkflow<ReceiveMessageWorkflow>()
.AddWorkflow<SendOneMessageWorkflow>()
.AddWorkflow<SendOneMessageWithCorrelationIdWorkflow>()
.ConfigureServices(services =>
{
services
.AddSingleton<ServiceBusClient>(_serviceBusClient)
.AddSingleton<IServiceBusProcessorManager, ServiceBusProcessorManager>()
.AddSingleton<IWorkerManager, WorkerManager>()
.AddSingleton(_resetEventManager)
.AddSingleton(sp =>
{
var options = sp.GetRequiredService<IOptions<MediatorOptions>>().Value;
return ActivatorUtilities.CreateInstance<BackgroundCommandSenderHostedService>(sp, options.CommandWorkerCount);
})
.AddSingleton(sp =>
{
var options = sp.GetRequiredService<IOptions<MediatorOptions>>().Value;
return ActivatorUtilities.CreateInstance<BackgroundEventPublisherHostedService>(sp, options.NotificationWorkerCount);
})
;
})
.Build();
_worker = _services.GetRequiredService<IWorkerManager>();
_hosted = _services.GetRequiredService<BackgroundCommandSenderHostedService>();
_backgroundEventService = _services.GetRequiredService<BackgroundEventPublisherHostedService>();
_sbProcessorManager = _services.GetRequiredService<IServiceBusProcessorManager>();
_testOutputHelper = testOutputHelper;
}
[Fact(DisplayName = "2 Receive - Sending 1 message - Should Block")]
public async Task Receive_1_Message_Should_Block_If_One_Receive()
{
_sbProcessorManager.Init("topicName", "subscriptionName");
_sbProcessorManager.Init("topicName1", "subscription1");
//Init waitEvent :
_resetEventManager.Init("receive1");
_resetEventManager.Init("receive2");
//Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
//Start Workflow
var workflowDefinitionId = typeof(ReceiveMessageWorkflow).Name;
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Start Worker to send Message on topicName/suscriptionName
await _worker.StartWorkerAsync("topicName", "subscriptionName");
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, null);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
//Wait for receiving second message
var wait2 = _resetEventManager.Get("receive2").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait2 : {wait2}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We don't send 2 messages so Workflow must be
* Running
* Suspended
*
* with a timeout on Wait for 2nd Message. (False ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Running, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, lastWorkflowState.SubStatus);
Assert.True(wait1);
Assert.False(wait2);
}
private async Task InitRegistryAndBackGroundServiceWorkerAsync()
{
//Init Registries to use StartWorkflow
await _services.PopulateRegistriesAsync();
//Start background services for CommandHandler
await _hosted.StartAsync(CancellationToken.None);
await _backgroundEventService.StartAsync(CancellationToken.None);
}
[Fact(DisplayName = "1 Receive - Sending 1 message - Should Finished")]
public async Task Receive_1_Message_Should_Finish_With_One_Receive()
{
_sbProcessorManager.Init("topicName", "subscriptionName");
//Init waitEvent :
_resetEventManager.Init("receive1");
//Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
//Start Workflow
var workflowDefinitionId = typeof(ReceiveOneMessageWorkflow).Name;
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Start Worker to send Message on topicName/suscriptionName
await _worker.StartWorkerAsync("topicName", "subscriptionName");
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, null);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We sent 1 message so Workflow must be
* Finished
* Finished
*
* with no timeout for the first message. (True ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Finished, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Finished, lastWorkflowState.SubStatus);
Assert.True(wait1);
}
[Fact(DisplayName = "1 Send - 1 Receive - Should Finished - without race condition")]
public async Task Send_1_Message_And_Receive_Response()
{
_sbProcessorManager.Init("topicName", "subscriptionName");
//Init waitEvent :
_resetEventManager.Init("receive1");
//Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
var senderMock = Substitute.For<ServiceBusSender>();
_serviceBusClient
.CreateSender(Arg.Any<string>())
.Returns(senderMock);
senderMock
.SendMessageAsync(Arg.Any<ServiceBusMessage>(), Arg.Any<CancellationToken>())
.ReturnsForAnyArgs(async (callback) =>
{
var sb = callback.Arg<ServiceBusMessage>();
var c = callback.Arg<CancellationToken>();
_testOutputHelper.WriteLine("Sending Message from activity");
await _worker.StartWorkerAsync("topicName", "subscriptionName");
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, null);
});
//Start Workflow
var workflowDefinitionId = typeof(SendOneMessageWorkflow).Name;
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We sent 1 message so Workflow must be
* Finished
* Finished
*
* with no timeout for the first message. (True ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Finished, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Finished, lastWorkflowState.SubStatus);
Assert.True(wait1);
}
[Fact(DisplayName = "1 Send - 1 Receive - with correlationId - Should Finished - without race condition")]
public async Task Send_1_Message_And_Correlate_Receive_Response()
{
var correlationId = "EEE3D9CC-2279-4CE5-8F4F-FC6C65BF8814";
_sbProcessorManager.Init("topicName", "subscriptionName");
//Init waitEvent :
_resetEventManager.Init("receive1");
//Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
var senderMock = Substitute.For<ServiceBusSender>();
_serviceBusClient
.CreateSender(Arg.Any<string>())
.Returns(senderMock);
senderMock
.SendMessageAsync(Arg.Any<ServiceBusMessage>(), Arg.Any<CancellationToken>())
.ReturnsForAnyArgs(async (callback) =>
{
var sb = callback.Arg<ServiceBusMessage>();
var c = callback.Arg<CancellationToken>();
_testOutputHelper.WriteLine("Sending Message from activity");
await _worker.StartWorkerAsync("topicName", "subscriptionName");
await _sbProcessorManager
.Get("topicName", "subscriptionName")
.SendMessage<dynamic>(new { hello = "world" }, correlationId);
});
//Start Workflow
var workflowDefinitionId = typeof(SendOneMessageWithCorrelationIdWorkflow).Name;
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
await Task.Delay(500); //Todo find how to remove delay
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We sent 1 message so Workflow must be
* Finished
* Finished
*
* with no timeout for the first message. (True ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Finished, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Finished, lastWorkflowState.SubStatus);
Assert.True(wait1);
}
[Fact(DisplayName = "1 Receive - Listening from other topic - Should not trigger")]
public async Task Receive_1_Message_Should_Not_Trigger()
{
_sbProcessorManager.Init("topicName1", "subscriptionName1");
//Init waitEvent :
_resetEventManager.Init("receive1");
_resetEventManager.Init("receive2");
//Init BackGround
await InitRegistryAndBackGroundServiceWorkerAsync();
//Start Workflow
var workflowDefinitionId = typeof(ReceiveMessageWorkflow).Name;
var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary<string, object>(), Common.Models.VersionOptions.Published);
var workflowRuntime = _services.GetRequiredService<IWorkflowRuntime>();
var workflowState = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, startWorkflowOptions);
/*
* Workflow don't receive any message so it should be
* Running
* Suspended
*/
Assert.Equal(WorkflowStatus.Running, workflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, workflowState.SubStatus);
//Start Worker to send Message on topicName/suscriptionName
await _worker.StartWorkerAsync("topicName1", "subscriptionName1");
await _sbProcessorManager
.Get("topicName1", "subscriptionName1")
.SendMessage<dynamic>(new { hello = "world" }, null);
//Wait for receiving first message
var wait1 = _resetEventManager.Get("receive1").WaitOne(TimeSpan.FromSeconds(5));
_testOutputHelper.WriteLine($"wait1 : {wait1}");
var lastWorkflowState = await workflowRuntime.ExportWorkflowStateAsync(workflowState.WorkflowInstanceId);
/*
* We don't send 2 messages so Workflow must be
* Running
* Suspended
*
* with a timeout on Wait for 2nd Message. (False ResetEvent)
*/
Assert.NotNull(lastWorkflowState);
Assert.Equal(WorkflowStatus.Running, lastWorkflowState.Status);
Assert.Equal(WorkflowSubStatus.Suspended, lastWorkflowState.SubStatus);
Assert.False(wait1);
}
public void Dispose()
{
_hosted.StopAsync(new CancellationToken()).GetAwaiter().GetResult();
}
}
}