159 lines
4.9 KiB
C#
159 lines
4.9 KiB
C#
using System.Threading;
|
|
using xCommons.Extensions;
|
|
using xEventService.Models;
|
|
using xExceptions.Constants;
|
|
using System.Threading.Tasks;
|
|
using xEventService.Constants;
|
|
using xEventService.Extensions;
|
|
using xEventService.Interfaces;
|
|
using xEventService.Configuration;
|
|
|
|
namespace xEventService.Providers
|
|
{
|
|
public class XEventServiceProvider : IXEventServiceProvider
|
|
{
|
|
//
|
|
private readonly XEventServiceConfiguration configuration;
|
|
private readonly IXKafkaProducerService kafkaProducerService;
|
|
private readonly IXZeroMQProducerService zeroMQProducerService;
|
|
private readonly IXRabbitMQProducerService rabbitMQProducerService;
|
|
|
|
public XEventServiceProvider(
|
|
XEventServiceConfiguration configuration,
|
|
IXKafkaProducerService kafkaProducerService = null,
|
|
IXZeroMQProducerService zeroMQProducerService = null,
|
|
IXRabbitMQProducerService rabbitMQProducerService = null
|
|
)
|
|
{
|
|
//
|
|
this.configuration = configuration;
|
|
this.kafkaProducerService = kafkaProducerService;
|
|
this.zeroMQProducerService = zeroMQProducerService;
|
|
this.rabbitMQProducerService = rabbitMQProducerService;
|
|
|
|
//
|
|
// Validate ...
|
|
var isValid = configuration.IsValid() &&
|
|
configuration.Broker == XEventBroker.XKafka
|
|
? kafkaProducerService != null
|
|
: configuration.Broker == XEventBroker.XZeroMQ
|
|
? zeroMQProducerService != null
|
|
: configuration.Broker == XEventBroker.XRabbitMQ
|
|
? rabbitMQProducerService != null
|
|
: false;
|
|
if (!isValid)
|
|
{
|
|
XException.InvalidConfiguration.Throw();
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Send Message to Broker Using XEventRequest instance ...
|
|
/// </summary>
|
|
/// <param name="request"></param>
|
|
/// <param name="cancellationToken"></param>
|
|
/// <returns></returns>
|
|
public async Task<bool> SendMessageAsync(
|
|
XEventRequest request,
|
|
CancellationToken cancellationToken = default
|
|
)
|
|
{
|
|
//
|
|
// Validate ...
|
|
if (!request.IsValid())
|
|
{
|
|
XException.InvalidArgs.Throw();
|
|
}
|
|
|
|
//
|
|
var result = await SendMessageAsync(
|
|
topic: request.Topic,
|
|
message: request.Message,
|
|
cancellationToken: cancellationToken
|
|
);
|
|
|
|
//
|
|
return result;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Send Message to Broker Using XMessage instance ...
|
|
/// </summary>
|
|
/// <param name="message"></param>
|
|
/// <param name="topic"></param>
|
|
/// <param name="cancellationToken"></param>
|
|
/// <returns></returns>
|
|
public async Task<bool> SendMessageAsync(
|
|
XEventMessage message,
|
|
string topic = XEventServiceConstants.XDefaultTopic,
|
|
CancellationToken cancellationToken = default
|
|
)
|
|
{
|
|
//
|
|
// Validate Args ...
|
|
if (!message.IsValid())
|
|
{
|
|
XException.InvalidArgs.Throw();
|
|
}
|
|
|
|
//
|
|
var result = false;
|
|
switch (configuration.Broker)
|
|
{
|
|
//
|
|
case XEventBroker.XKafka:
|
|
result = await kafkaProducerService.SendMessageAsync(
|
|
topic: topic,
|
|
message: message,
|
|
cancellationToken: cancellationToken
|
|
);
|
|
break;
|
|
|
|
//
|
|
case XEventBroker.XZeroMQ:
|
|
result = await zeroMQProducerService.SendMessageAsync(
|
|
topic: topic,
|
|
message: message,
|
|
cancellationToken: cancellationToken
|
|
);
|
|
break;
|
|
|
|
//
|
|
case XEventBroker.XRabbitMQ:
|
|
result = await rabbitMQProducerService.SendMessageAsync(
|
|
topic: topic,
|
|
message: message,
|
|
cancellationToken: cancellationToken
|
|
);
|
|
break;
|
|
}
|
|
|
|
//
|
|
return result;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Dispose ...
|
|
/// </summary>
|
|
public void Dispose()
|
|
{
|
|
//
|
|
if (kafkaProducerService != null)
|
|
{
|
|
kafkaProducerService.Dispose();
|
|
}
|
|
|
|
//
|
|
if (zeroMQProducerService != null)
|
|
{
|
|
zeroMQProducerService.Dispose();
|
|
}
|
|
|
|
//
|
|
if (rabbitMQProducerService != null)
|
|
{
|
|
rabbitMQProducerService.Dispose();
|
|
}
|
|
}
|
|
}
|
|
} |