add supports for ZeroMQ, XRabbitMQ, XKafka ...
This commit is contained in:
@@ -0,0 +1,198 @@
|
||||
using System;
|
||||
using System.Text;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using NetMQ;
|
||||
using NetMQ.Sockets;
|
||||
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 XZeroMQProducerService : XBaseEventProducerService, IXZeroMQProducerService
|
||||
{
|
||||
//
|
||||
private NetMQSocket socket;
|
||||
private readonly object socketLock = new();
|
||||
|
||||
public XZeroMQProducerService(
|
||||
XEventServiceConfiguration configuration,
|
||||
ILogger<XZeroMQProducerService> logger
|
||||
) : base(configuration, logger)
|
||||
{
|
||||
//
|
||||
// Validate Configurations ...
|
||||
if (!configuration.IsValid() ||
|
||||
!configuration.ZeroMQ.IsValid() ||
|
||||
configuration.Broker != XEventBroker.XZeroMQ
|
||||
)
|
||||
{
|
||||
XException.InvalidConfiguration.Throw();
|
||||
}
|
||||
|
||||
//
|
||||
InitializeSocket();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Sending Message with Full Request Descriptor ...
|
||||
/// </summary>
|
||||
/// <param name="request"></param>
|
||||
/// <param name="cancellationToken"></param>
|
||||
/// <returns></returns>
|
||||
public async Task<bool> SendMessageAsync(
|
||||
XEventRequest request,
|
||||
CancellationToken cancellationToken = default
|
||||
)
|
||||
{
|
||||
//
|
||||
var result = request.IsValid();
|
||||
if (!result)
|
||||
{
|
||||
XException.InvalidArgs.Throw();
|
||||
}
|
||||
|
||||
//
|
||||
result = await SendMessageAsync(
|
||||
message: request.Message,
|
||||
topic: request.Topic,
|
||||
cancellationToken: cancellationToken
|
||||
);
|
||||
|
||||
//
|
||||
return result;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Sending a Message to Topic ...
|
||||
/// </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
|
||||
)
|
||||
{
|
||||
//
|
||||
var result = message.IsValid();
|
||||
if (!result)
|
||||
{
|
||||
XException.InvalidArgs.Throw();
|
||||
}
|
||||
|
||||
//
|
||||
// Normalize Topic ...
|
||||
if (topic.IsNullOrEmpty())
|
||||
{
|
||||
topic = XEventServiceConstants.XDefaultTopic;
|
||||
}
|
||||
|
||||
//
|
||||
try
|
||||
{
|
||||
//
|
||||
var json = message.ToJSON(camelCase: true);
|
||||
var body = Encoding.UTF8.GetBytes(json);
|
||||
await Task.Run(() =>
|
||||
{
|
||||
//
|
||||
lock (socketLock)
|
||||
{
|
||||
//
|
||||
if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
|
||||
{
|
||||
socket.SendMoreFrame(topic).SendFrame(body);
|
||||
}
|
||||
else
|
||||
{
|
||||
socket.SendFrame(body);
|
||||
}
|
||||
}
|
||||
}, cancellationToken);
|
||||
|
||||
//
|
||||
result = true;
|
||||
}
|
||||
catch (Exception)
|
||||
{
|
||||
result = false;
|
||||
}
|
||||
|
||||
//
|
||||
return result;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Dispose Object ...
|
||||
/// </summary>
|
||||
public void Dispose()
|
||||
{
|
||||
//
|
||||
try
|
||||
{
|
||||
//
|
||||
lock (socketLock)
|
||||
{
|
||||
//
|
||||
if (socket != null)
|
||||
{
|
||||
//
|
||||
socket.Close();
|
||||
socket.Dispose();
|
||||
socket = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
catch { }
|
||||
|
||||
//
|
||||
GC.SuppressFinalize(this);
|
||||
}
|
||||
|
||||
//
|
||||
#region Private ...
|
||||
/// <summary>
|
||||
/// Initialize Socket ...
|
||||
/// </summary>
|
||||
private void InitializeSocket()
|
||||
{
|
||||
//
|
||||
lock (socketLock)
|
||||
{
|
||||
//
|
||||
if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
|
||||
{
|
||||
//
|
||||
socket = new PublisherSocket();
|
||||
socket.Options.SendHighWatermark = configuration.ZeroMQ.SendHighWaterMark;
|
||||
}
|
||||
else
|
||||
{
|
||||
//
|
||||
socket = new PushSocket();
|
||||
socket.Options.SendHighWatermark = configuration.ZeroMQ.SendHighWaterMark;
|
||||
}
|
||||
|
||||
//
|
||||
if (configuration.ZeroMQ.IsServer)
|
||||
{
|
||||
socket.Bind(configuration.Url);
|
||||
}
|
||||
else
|
||||
{
|
||||
socket.Connect(configuration.Url);
|
||||
}
|
||||
}
|
||||
}
|
||||
#endregion
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user