diff --git a/Elsa.sln b/Elsa.sln index 23dd0aee3..4d05d3ff6 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -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 diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Contracts/IServiceBusProcessorManager.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Contracts/IServiceBusProcessorManager.cs similarity index 66% rename from test/integration/Elsa.ServiceBusIntegrationTests/Contracts/IServiceBusProcessorManager.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Contracts/IServiceBusProcessorManager.cs index 69cac4d43..7664f5e60 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Contracts/IServiceBusProcessorManager.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Contracts/IServiceBusProcessorManager.cs @@ -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 { diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Contracts/ITestResetEventManager.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Contracts/ITestResetEventManager.cs similarity index 76% rename from test/integration/Elsa.ServiceBusIntegrationTests/Contracts/ITestResetEventManager.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Contracts/ITestResetEventManager.cs index aed3a0979..ee9d41d7c 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Contracts/ITestResetEventManager.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Contracts/ITestResetEventManager.cs @@ -1,4 +1,4 @@ -namespace Elsa.ServiceBusIntegrationTests.Contracts +namespace Elsa.ServiceBus.IntegrationTests.Contracts { public interface ITestResetEventManager { diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Elsa.ServiceBusIntegrationTests.csproj b/test/integration/Elsa.ServiceBus.IntegrationTests/Elsa.ServiceBus.IntegrationTests.csproj similarity index 100% rename from test/integration/Elsa.ServiceBusIntegrationTests/Elsa.ServiceBusIntegrationTests.csproj rename to test/integration/Elsa.ServiceBus.IntegrationTests/Elsa.ServiceBus.IntegrationTests.csproj diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/GlobalUsings.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/GlobalUsings.cs similarity index 100% rename from test/integration/Elsa.ServiceBusIntegrationTests/GlobalUsings.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/GlobalUsings.cs diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Helpers/ServiceBusProcessorManager.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/ServiceBusProcessorManager.cs similarity index 79% rename from test/integration/Elsa.ServiceBusIntegrationTests/Helpers/ServiceBusProcessorManager.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/ServiceBusProcessorManager.cs index 6116ad176..1c2e43517 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Helpers/ServiceBusProcessorManager.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/ServiceBusProcessorManager.cs @@ -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(s => s == topic), Arg.Any(), Arg.Any() - ) - .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); diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Helpers/ServiceBusProcessorTest.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/ServiceBusProcessorTest.cs similarity index 76% rename from test/integration/Elsa.ServiceBusIntegrationTests/Helpers/ServiceBusProcessorTest.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/ServiceBusProcessorTest.cs index a13f192e0..35efe3233 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Helpers/ServiceBusProcessorTest.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/ServiceBusProcessorTest.cs @@ -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 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 payload, string correlationId, int deliveryCount = 1) { var payloadJson = JsonSerializer.Serialize(payload); var props = new Dictionary() { }; @@ -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"); - } } } \ No newline at end of file diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Helpers/TestResetEventManager.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/TestResetEventManager.cs similarity index 87% rename from test/integration/Elsa.ServiceBusIntegrationTests/Helpers/TestResetEventManager.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/TestResetEventManager.cs index ac77e23d2..833569264 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Helpers/TestResetEventManager.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Helpers/TestResetEventManager.cs @@ -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 { diff --git a/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/ServiceBus/ServiceBusTests.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/ServiceBus/ServiceBusTests.cs new file mode 100644 index 000000000..04102031f --- /dev/null +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/ServiceBus/ServiceBusTests.cs @@ -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(); + 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() + .AddWorkflow() + .AddWorkflow() + .AddWorkflow() + .ConfigureServices(services => + { + services + .AddSingleton(_serviceBusClient) + .AddSingleton() + .AddSingleton() + .AddSingleton(_resetEventManager) + + .AddSingleton(sp => + { + var options = sp.GetRequiredService>().Value; + return ActivatorUtilities.CreateInstance(sp, options.CommandWorkerCount); + }) + .AddSingleton(sp => + { + var options = sp.GetRequiredService>().Value; + return ActivatorUtilities.CreateInstance(sp, options.NotificationWorkerCount); + }) + ; + }) + .Build(); + + _worker = _services.GetRequiredService(); + _hosted = _services.GetRequiredService(); + _backgroundEventService = _services.GetRequiredService(); + _sbProcessorManager = _services.GetRequiredService(); + + _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(), Common.Models.VersionOptions.Published); + var workflowRuntime = _services.GetRequiredService(); + 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(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(), Common.Models.VersionOptions.Published); + var workflowRuntime = _services.GetRequiredService(); + 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(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(); + _serviceBusClient + .CreateSender(Arg.Any()) + .Returns(senderMock); + + senderMock + .SendMessageAsync(Arg.Any(), Arg.Any()) + .ReturnsForAnyArgs(async (callback) => + { + var sb = callback.Arg(); + var c = callback.Arg(); + + _testOutputHelper.WriteLine("Sending Message from activity"); + + await _worker.StartWorkerAsync("topicName", "subscriptionName", c); + await _sbProcessorManager + .Get("topicName", "subscriptionName") + .SendMessage(new { hello = "world" }, null!); + }); + + + // Start Workflow + const string workflowDefinitionId = nameof(SendOneMessageWorkflow); + var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary(), Common.Models.VersionOptions.Published); + var workflowRuntime = _services.GetRequiredService(); + 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(); + _serviceBusClient + .CreateSender(Arg.Any()) + .Returns(senderMock); + + senderMock + .SendMessageAsync(Arg.Any(), Arg.Any()) + .ReturnsForAnyArgs(async (callback) => + { + var sb = callback.Arg(); + var c = callback.Arg(); + + _testOutputHelper.WriteLine("Sending Message from activity"); + + await _worker.StartWorkerAsync("topicName", "subscriptionName"); + await _sbProcessorManager + .Get("topicName", "subscriptionName") + .SendMessage(new { hello = "world" }, correlationId); + }); + + //Start Workflow + var workflowDefinitionId = nameof(SendOneMessageWithCorrelationIdWorkflow); + var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary(), Common.Models.VersionOptions.Published); + var workflowRuntime = _services.GetRequiredService(); + 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(), Common.Models.VersionOptions.Published); + var workflowRuntime = _services.GetRequiredService(); + 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(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(); + } +} \ No newline at end of file diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/ReceiveMessageWorkflow.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/ReceiveMessageWorkflow.cs similarity index 87% rename from test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/ReceiveMessageWorkflow.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/ReceiveMessageWorkflow.cs index 7806af36c..596de96e1 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/ReceiveMessageWorkflow.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/ReceiveMessageWorkflow.cs @@ -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 { diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/ReceiveOneMessageWorkflow.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/ReceiveOneMessageWorkflow.cs similarity index 85% rename from test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/ReceiveOneMessageWorkflow.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/ReceiveOneMessageWorkflow.cs index e29f15f0a..5b2214602 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/ReceiveOneMessageWorkflow.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/ReceiveOneMessageWorkflow.cs @@ -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 { diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/SendOneMessageWorkflow.cs b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/SendOneMessageWorkflow.cs similarity index 92% rename from test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/SendOneMessageWorkflow.cs rename to test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/SendOneMessageWorkflow.cs index dc9f293d8..191df9b42 100644 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/Workflows/SendOneMessageWorkflow.cs +++ b/test/integration/Elsa.ServiceBus.IntegrationTests/Scenarios/Workflows/SendOneMessageWorkflow.cs @@ -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 { diff --git a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/ServiceBus/ServiceBusTests.cs b/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/ServiceBus/ServiceBusTests.cs deleted file mode 100644 index fda26c927..000000000 --- a/test/integration/Elsa.ServiceBusIntegrationTests/Scenarios/ServiceBus/ServiceBusTests.cs +++ /dev/null @@ -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(); - - 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() - .AddWorkflow() - .AddWorkflow() - .AddWorkflow() - .ConfigureServices(services => - { - services - .AddSingleton(_serviceBusClient) - .AddSingleton() - - .AddSingleton() - .AddSingleton(_resetEventManager) - - .AddSingleton(sp => - { - var options = sp.GetRequiredService>().Value; - return ActivatorUtilities.CreateInstance(sp, options.CommandWorkerCount); - }) - .AddSingleton(sp => - { - var options = sp.GetRequiredService>().Value; - return ActivatorUtilities.CreateInstance(sp, options.NotificationWorkerCount); - }) - ; - }) - .Build(); - - _worker = _services.GetRequiredService(); - _hosted = _services.GetRequiredService(); - _backgroundEventService = _services.GetRequiredService(); - _sbProcessorManager = _services.GetRequiredService(); - - _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(), Common.Models.VersionOptions.Published); - var workflowRuntime = _services.GetRequiredService(); - 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(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(), Common.Models.VersionOptions.Published); - var workflowRuntime = _services.GetRequiredService(); - 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(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(); - _serviceBusClient - .CreateSender(Arg.Any()) - .Returns(senderMock); - - - senderMock - .SendMessageAsync(Arg.Any(), Arg.Any()) - .ReturnsForAnyArgs(async (callback) => - { - var sb = callback.Arg(); - var c = callback.Arg(); - - _testOutputHelper.WriteLine("Sending Message from activity"); - - await _worker.StartWorkerAsync("topicName", "subscriptionName"); - await _sbProcessorManager - .Get("topicName", "subscriptionName") - .SendMessage(new { hello = "world" }, null); - }); - - - //Start Workflow - var workflowDefinitionId = typeof(SendOneMessageWorkflow).Name; - var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary(), Common.Models.VersionOptions.Published); - var workflowRuntime = _services.GetRequiredService(); - 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(); - _serviceBusClient - .CreateSender(Arg.Any()) - .Returns(senderMock); - - senderMock - .SendMessageAsync(Arg.Any(), Arg.Any()) - .ReturnsForAnyArgs(async (callback) => - { - var sb = callback.Arg(); - var c = callback.Arg(); - - _testOutputHelper.WriteLine("Sending Message from activity"); - - await _worker.StartWorkerAsync("topicName", "subscriptionName"); - await _sbProcessorManager - .Get("topicName", "subscriptionName") - .SendMessage(new { hello = "world" }, correlationId); - }); - - //Start Workflow - var workflowDefinitionId = typeof(SendOneMessageWithCorrelationIdWorkflow).Name; - var startWorkflowOptions = new StartWorkflowRuntimeOptions(null, new Dictionary(), Common.Models.VersionOptions.Published); - var workflowRuntime = _services.GetRequiredService(); - 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(), Common.Models.VersionOptions.Published); - var workflowRuntime = _services.GetRequiredService(); - 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(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(); - } - } -} \ No newline at end of file