2023-10-19 18:00:43 +00:00
|
|
|
|
using Azure.Messaging.ServiceBus;
|
2023-10-20 11:38:40 +00:00
|
|
|
|
using Elsa.ServiceBus.IntegrationTests.Contracts;
|
2023-10-19 18:00:43 +00:00
|
|
|
|
using NSubstitute;
|
|
|
|
|
|
using Xunit.Abstractions;
|
|
|
|
|
|
|
2023-12-18 19:20:39 +00:00
|
|
|
|
namespace Elsa.ServiceBus.IntegrationTests.Helpers;
|
2023-10-20 11:38:40 +00:00
|
|
|
|
|
2023-12-18 19:20:39 +00:00
|
|
|
|
public class ServiceBusProcessorManager(ServiceBusClient serviceBusClient, ITestOutputHelper testOutputHelper) : IServiceBusProcessorManager
|
|
|
|
|
|
{
|
|
|
|
|
|
private readonly IDictionary<(string topicName, string subscriptionName), ServiceBusProcessorTest> _serviceBusProcessors = new Dictionary<(string topicName, string subscriptionName), ServiceBusProcessorTest>();
|
2023-10-19 18:00:43 +00:00
|
|
|
|
|
2023-12-18 19:20:39 +00:00
|
|
|
|
public ServiceBusProcessorTest Init(string topic, string subscription)
|
|
|
|
|
|
{
|
|
|
|
|
|
_serviceBusProcessors.TryGetValue((topic, subscription), out var processor);
|
|
|
|
|
|
if (processor != null)
|
|
|
|
|
|
throw new ArgumentException($"ServiceBusProcessor with topic:{topic}, subscription:{subscription} already initialized");
|
2023-10-19 18:00:43 +00:00
|
|
|
|
|
2023-12-18 19:20:39 +00:00
|
|
|
|
processor = new ServiceBusProcessorTest(testOutputHelper);
|
|
|
|
|
|
serviceBusClient
|
|
|
|
|
|
.CreateProcessor(Arg.Is<string>(s => s == topic),
|
|
|
|
|
|
Arg.Any<string>(),
|
|
|
|
|
|
Arg.Any<ServiceBusProcessorOptions>()
|
|
|
|
|
|
)
|
|
|
|
|
|
.Returns((c) =>
|
|
|
|
|
|
{
|
|
|
|
|
|
testOutputHelper.WriteLine("ServiceBusClient");
|
|
|
|
|
|
return processor;
|
|
|
|
|
|
});
|
2023-10-19 18:00:43 +00:00
|
|
|
|
|
2023-12-18 19:20:39 +00:00
|
|
|
|
_serviceBusProcessors.Add((topic, subscription), processor);
|
|
|
|
|
|
return processor;
|
|
|
|
|
|
}
|
2023-10-20 11:38:40 +00:00
|
|
|
|
|
2023-12-18 19:20:39 +00:00
|
|
|
|
public ServiceBusProcessorTest Get(string topic, string subscription)
|
|
|
|
|
|
{
|
|
|
|
|
|
_serviceBusProcessors.TryGetValue((topic, subscription), out var processor);
|
|
|
|
|
|
if (processor == null)
|
|
|
|
|
|
throw new ArgumentException($"ServiceBusProcessor with topic:{topic}, subscription:{subscription} not exist");
|
2023-10-19 18:00:43 +00:00
|
|
|
|
|
2023-12-18 19:20:39 +00:00
|
|
|
|
return processor;
|
2023-10-19 18:00:43 +00:00
|
|
|
|
}
|
|
|
|
|
|
}
|