Complete xEventService ...
This commit is contained in:
@@ -0,0 +1,159 @@
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using xCommons.Extensions;
|
||||
using xEventService.Configuration;
|
||||
using xEventService.Constants;
|
||||
using xEventService.Extensions;
|
||||
using xEventService.Interfaces;
|
||||
using xEventService.Models;
|
||||
using xExceptions.Constants;
|
||||
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user