diff --git a/Configuration/XEventServiceConfiguration.cs b/Configuration/XEventServiceConfiguration.cs
index f291fb0..6409b48 100644
--- a/Configuration/XEventServiceConfiguration.cs
+++ b/Configuration/XEventServiceConfiguration.cs
@@ -23,7 +23,13 @@ namespace xEventService.Configuration
/// Allowed Message Broker ...
///
public XEventBroker Broker { get; set; } = XEventBroker.None;
-
+
+ ///
+ /// ZeroMQ Specific Options ...
+ /// Only used when Broker is XZeroMQ ...
+ ///
+ public XZeroMQOptions ZeroMQ { get; set; } = new XZeroMQOptions();
+
///
/// RabbitMQ Specific Options ...
/// Only used when Broker is XRabbitMQ ...
diff --git a/Interfaces/IXRabbitMQProducerService.cs b/Interfaces/IXRabbitMQProducerService.cs
index 46bd0f4..a8794f4 100644
--- a/Interfaces/IXRabbitMQProducerService.cs
+++ b/Interfaces/IXRabbitMQProducerService.cs
@@ -1,4 +1,3 @@
-using System;
using System.Threading;
using System.Threading.Tasks;
using xEventService.Models;
@@ -29,6 +28,6 @@ namespace xEventService.Interfaces
string exchange = null,
string routingKey = null,
CancellationToken cancellationToken = default
- );
+ );
}
}
\ No newline at end of file
diff --git a/Interfaces/IXZeroMQProducerService.cs b/Interfaces/IXZeroMQProducerService.cs
index 894d934..f89b7dc 100644
--- a/Interfaces/IXZeroMQProducerService.cs
+++ b/Interfaces/IXZeroMQProducerService.cs
@@ -1,7 +1,5 @@
namespace xEventService.Interfaces
{
public interface IXZeroMQProducerService : IXProducerServiceBase
- {
-
- }
+ { }
}
\ No newline at end of file
diff --git a/Providers/XBaseEventConsumerBackgroundService.cs b/Providers/XBaseEventConsumerBackgroundService.cs
index 7dc94e5..f3051e7 100644
--- a/Providers/XBaseEventConsumerBackgroundService.cs
+++ b/Providers/XBaseEventConsumerBackgroundService.cs
@@ -22,12 +22,12 @@ namespace xEventService.Providers
///
/// Allowed Senders ...
///
- public readonly string[] allowedSenders;
+ protected readonly string[] allowedSenders;
///
/// Allowed Actions ...
///
- public readonly string[] allowedActions;
+ protected readonly string[] allowedActions;
///
/// Configuration ...
@@ -41,10 +41,7 @@ namespace xEventService.Providers
XEventServiceConfiguration configuration,
ILogger logger,
string topic = XEventServiceConstants.XDefaultTopic
- ) : base(
- logger: logger,
- serviceProvider: serviceProvider
- )
+ ) : base(serviceProvider, logger)
{
//
// Validate Actions ...
diff --git a/Providers/XBaseEventProducerService.cs b/Providers/XBaseEventProducerService.cs
new file mode 100644
index 0000000..2004bed
--- /dev/null
+++ b/Providers/XBaseEventProducerService.cs
@@ -0,0 +1,31 @@
+using Microsoft.Extensions.Logging;
+using xCommons.Extensions;
+using xEventService.Configuration;
+using xEventService.Extensions;
+using xExceptions.Constants;
+
+namespace xEventService.Providers
+{
+ public abstract class XBaseEventProducerService
+ {
+ protected readonly XEventServiceConfiguration configuration;
+ protected readonly ILogger logger;
+
+ protected XBaseEventProducerService(
+ XEventServiceConfiguration configuration,
+ ILogger logger
+ )
+ {
+ //
+ this.logger = logger;
+ this.configuration = configuration;
+
+ //
+ // Validate Configurations ...
+ if (!configuration.IsValid())
+ {
+ XException.InvalidConfiguration.Throw();
+ }
+ }
+ }
+}
\ No newline at end of file
diff --git a/Providers/XKafkaConsumerServiceBase.cs b/Providers/XKafkaConsumerServiceBase.cs
index a55d452..62fd8c6 100644
--- a/Providers/XKafkaConsumerServiceBase.cs
+++ b/Providers/XKafkaConsumerServiceBase.cs
@@ -6,6 +6,8 @@ using Microsoft.Extensions.Logging;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
+using xEventService.Extensions;
+using xEventService.Models;
using xExceptions.Constants;
namespace xEventService.Providers
@@ -87,10 +89,21 @@ namespace xEventService.Providers
)
{
//
- await DoWorkAsync(
- payload: consumeResult.Message,
- cancellationToken: cancellationToken
- );
+ var message = consumeResult
+ .Message
+ .Value
+ .FromJSON();
+ if (message.IsValid(
+ allowedSenders: allowedSenders,
+ allowedActions: allowedActions)
+ )
+ {
+ //
+ await DoWorkAsync(
+ payload: consumeResult.Message,
+ cancellationToken: cancellationToken
+ );
+ }
}
}
catch (ConsumeException)
@@ -105,13 +118,18 @@ namespace xEventService.Providers
}
///
- /// Dispose Implementation ...
+ /// Dispose Object ...
///
public override void Dispose()
{
//
- consumer.Close();
- consumer.Dispose();
+ if (consumer != null)
+ {
+ consumer.Close();
+ consumer.Dispose();
+ }
+
+ //
GC.SuppressFinalize(this);
base.Dispose();
}
diff --git a/Providers/XKafkaProducerService.cs b/Providers/XKafkaProducerService.cs
index 0c693fb..caccc35 100644
--- a/Providers/XKafkaProducerService.cs
+++ b/Providers/XKafkaProducerService.cs
@@ -13,21 +13,15 @@ using System;
namespace xEventService.Providers
{
- public class XKafkaProducerService : IXKafkaProducerService
+ public class XKafkaProducerService : XBaseEventProducerService, IXKafkaProducerService
{
private readonly IProducer producer;
- private readonly ILogger logger;
- private readonly XEventServiceConfiguration configuration;
public XKafkaProducerService(
ILogger logger,
XEventServiceConfiguration configuration
- )
+ ) : base(configuration, logger)
{
- //
- this.logger = logger;
- this.configuration = configuration;
-
//
// Validate Configurations ...
if (!configuration.IsValid() ||
@@ -132,7 +126,12 @@ namespace xEventService.Providers
public void Dispose()
{
//
- producer.Dispose();
+ if (producer != null)
+ {
+ producer.Dispose();
+ }
+
+ //
GC.SuppressFinalize(this);
}
}
diff --git a/Providers/XRabbitMQConsumerServiceBase.cs b/Providers/XRabbitMQConsumerServiceBase.cs
index 37ad228..5d27e10 100644
--- a/Providers/XRabbitMQConsumerServiceBase.cs
+++ b/Providers/XRabbitMQConsumerServiceBase.cs
@@ -19,6 +19,7 @@ namespace xEventService.Providers
///
public abstract class XRabbitMQConsumerServiceBase : XBaseEventConsumerBackgroundService
{
+ //
protected IModel channel;
protected IConnection connection;
protected AsyncEventingBasicConsumer consumer;
@@ -79,7 +80,11 @@ namespace xEventService.Providers
var body = ea.Body.ToArray();
var json = Encoding.UTF8.GetString(body);
var message = json.FromJSON();
- if (!message.IsNullOrDefault())
+ if (!message.IsNullOrDefault() &&
+ message.IsValid(
+ allowedSenders: allowedSenders,
+ allowedActions: allowedActions)
+ )
{
//
await DoWorkAsync(
@@ -163,8 +168,18 @@ namespace xEventService.Providers
public override void Dispose()
{
//
- channel.Dispose();
- connection.Dispose();
+ if (channel != null)
+ {
+ channel.Dispose();
+ }
+
+ //
+ if (connection != null)
+ {
+ connection.Dispose();
+ }
+
+ //
GC.SuppressFinalize(this);
base.Dispose();
}
diff --git a/Providers/XRabbitMQProducerService.cs b/Providers/XRabbitMQProducerService.cs
index 3d27b63..b462a28 100644
--- a/Providers/XRabbitMQProducerService.cs
+++ b/Providers/XRabbitMQProducerService.cs
@@ -17,23 +17,17 @@ namespace xEventService.Providers
///
/// RabbitMQ Producer Service Implementation ...
///
- public class XRabbitMQProducerService : IXRabbitMQProducerService
+ public class XRabbitMQProducerService : XBaseEventProducerService, IXRabbitMQProducerService
{
private IModel channel;
private IConnection connection;
private readonly object connectionLock = new();
- private readonly ILogger logger;
- private readonly XEventServiceConfiguration configuration;
public XRabbitMQProducerService(
ILogger logger,
XEventServiceConfiguration configuration
- )
+ ) : base(configuration, logger)
{
- //
- this.logger = logger;
- this.configuration = configuration;
-
//
// Validate Configurations ...
if (!configuration.IsValid() ||
@@ -52,6 +46,9 @@ namespace xEventService.Providers
///
/// Sending Message with Full Request Descriptor ...
///
+ ///
+ ///
+ ///
public async Task SendMessageAsync(
XEventRequest request,
CancellationToken cancellationToken = default
@@ -77,9 +74,41 @@ namespace xEventService.Providers
return result;
}
+ ///
+ /// Sending a Message to Topic ...
+ ///
+ ///
+ ///
+ ///
+ ///
+ public async Task SendMessageAsync(
+ XEventMessage message,
+ string topic = XEventServiceConstants.XDefaultTopic,
+ CancellationToken cancellationToken = default
+ )
+ {
+ //
+ var result = await SendMessageAsync(
+ request: new XEventRequest
+ {
+ Topic = topic,
+ Message = message
+ },
+ cancellationToken: cancellationToken
+ );
+
+ //
+ return result;
+ }
+
///
/// Sending a Message to Exchange ...
///
+ ///
+ ///
+ ///
+ ///
+ ///
public async Task SendMessageAsync(
XEventMessage message,
string exchange = null,
@@ -114,7 +143,7 @@ namespace xEventService.Providers
{
//
// Ensure Connection is Alive ...
- await EnsureChannelAsync(cancellationToken);
+ EnsureChannelAsync();
//
// Serialize Message ...
@@ -170,14 +199,22 @@ namespace xEventService.Providers
try
{
//
- channel.Close();
- channel.Dispose();
- channel = null;
+ if (channel != null)
+ {
+ //
+ channel.Close();
+ channel.Dispose();
+ channel = null;
+ }
//
- connection.Close();
- connection.Dispose();
- connection = null;
+ if (connection != null)
+ {
+ //
+ connection.Close();
+ connection.Dispose();
+ connection = null;
+ }
}
catch
{ }
@@ -228,29 +265,10 @@ namespace xEventService.Providers
return result;
}
- ///
- /// Ensure Connection Established ...
- ///
- private async Task EnsureConnectionAsync()
- {
- //
- try
- {
- await EnsureChannelAsync();
- return true;
- }
- catch
- {
- return false;
- }
- }
-
///
/// Ensure Channel is Open and Exchange/Queue Declared ...
///
- private async Task EnsureChannelAsync(
- CancellationToken cancellationToken = default
- )
+ private void EnsureChannelAsync()
{
//
if (IsConnected())
diff --git a/Providers/XZeroMQConsumerServiceBase.cs b/Providers/XZeroMQConsumerServiceBase.cs
new file mode 100644
index 0000000..144563a
--- /dev/null
+++ b/Providers/XZeroMQConsumerServiceBase.cs
@@ -0,0 +1,187 @@
+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.Models;
+using xExceptions.Constants;
+
+namespace xEventService.Providers
+{
+ ///
+ /// an Abstract ZeroMQ Consumer Service ...
+ ///
+ public abstract class XZeroMQConsumerServiceBase : XBaseEventConsumerBackgroundService
+ {
+ //
+ private NetMQSocket socket;
+
+ protected XZeroMQConsumerServiceBase(
+ string[] allowedActions,
+ string[] allowedSenders,
+ IServiceProvider serviceProvider,
+ XEventServiceConfiguration configuration,
+ ILogger logger,
+ string topic = XEventServiceConstants.XDefaultTopic
+ ) : base(
+ topic: topic,
+ logger: logger,
+ configuration: configuration,
+ allowedActions: allowedActions,
+ allowedSenders: allowedSenders,
+ serviceProvider: serviceProvider
+ )
+ {
+ //
+ // Validating Configuration ...
+ if (configuration.Broker != XEventBroker.XKafka)
+ {
+ XException.InvalidConfiguration.Throw();
+ }
+
+ //
+ // Initialize ...
+ InitializeSocket();
+ }
+
+ ///
+ /// Default Action Execution for Background Services ...
+ ///
+ ///
+ ///
+ protected override async Task ExecuteAsync(
+ CancellationToken cancellationToken = default
+ )
+ {
+ //
+ await Task.Yield();
+ while (!cancellationToken.IsCancellationRequested)
+ {
+ //
+ try
+ {
+ //
+ string receivedTopic = string.Empty;
+ byte[] body = null;
+
+ if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
+ {
+ //
+ if (socket.TryReceiveFrameString(TimeSpan.FromMilliseconds(100), out receivedTopic))
+ {
+ //
+ if (socket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(100), out body))
+ {
+ // Process
+ }
+ }
+ }
+ else
+ {
+ if (socket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(100), out body))
+ {
+ // Process
+ }
+ }
+
+ if (body != null && body.Length > 0)
+ {
+ //
+ var json = Encoding.UTF8.GetString(body);
+ var message = json.FromJSON();
+ if (!message.IsNullOrDefault() &&
+ message.IsValid(
+ allowedSenders: allowedSenders,
+ allowedActions: allowedActions)
+ )
+ {
+ //
+ await DoWorkAsync(
+ payload: message,
+ cancellationToken: cancellationToken
+ );
+ }
+ }
+ else
+ {
+ await Task.Delay(10, cancellationToken);
+ }
+ }
+ catch (Exception)
+ {
+ await Task.Delay(1000, cancellationToken);
+ }
+ }
+ }
+
+ public override async Task StopAsync(CancellationToken cancellationToken)
+ {
+ //
+ socket.Close();
+ await base.StopAsync(cancellationToken);
+ }
+
+ ///
+ /// Dispose Object ...
+ ///
+ public override void Dispose()
+ {
+ //
+ if (socket != null)
+ {
+ socket.Dispose();
+ }
+
+ //
+ GC.SuppressFinalize(this);
+ base.Dispose();
+ }
+
+ //
+ #region Private ...
+ private void InitializeSocket()
+ {
+ //
+ if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
+ {
+ //
+ socket = new SubscriberSocket();
+ socket.Options.ReceiveHighWatermark = configuration.ZeroMQ.ReceiveHighWaterMark;
+
+ //
+ // Subscribe to specific topic or all
+ if (topic.IsNullOrEmpty() || topic == XEventServiceConstants.XWildcard)
+ {
+ ((SubscriberSocket)socket).SubscribeToAnyTopic();
+ }
+ else
+ {
+ ((SubscriberSocket)socket).Subscribe(topic);
+ }
+ }
+ else
+ {
+ //
+ socket = new PullSocket();
+ socket.Options.ReceiveHighWatermark = configuration.ZeroMQ.ReceiveHighWaterMark;
+ }
+
+ //
+ if (configuration.ZeroMQ.IsServer)
+ {
+ socket.Bind(configuration.Url);
+ }
+ else
+ {
+ socket.Connect(configuration.Url);
+ }
+ }
+ #endregion
+ }
+}
\ No newline at end of file
diff --git a/Providers/XZeroMQProducerService.cs b/Providers/XZeroMQProducerService.cs
new file mode 100644
index 0000000..c38c44d
--- /dev/null
+++ b/Providers/XZeroMQProducerService.cs
@@ -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 logger
+ ) : base(configuration, logger)
+ {
+ //
+ // Validate Configurations ...
+ if (!configuration.IsValid() ||
+ !configuration.ZeroMQ.IsValid() ||
+ configuration.Broker != XEventBroker.XZeroMQ
+ )
+ {
+ XException.InvalidConfiguration.Throw();
+ }
+
+ //
+ InitializeSocket();
+ }
+
+ ///
+ /// Sending Message with Full Request Descriptor ...
+ ///
+ ///
+ ///
+ ///
+ public async Task 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;
+ }
+
+ ///
+ /// Sending a Message to Topic ...
+ ///
+ ///
+ ///
+ ///
+ ///
+ public async Task 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;
+ }
+
+ ///
+ /// Dispose Object ...
+ ///
+ public void Dispose()
+ {
+ //
+ try
+ {
+ //
+ lock (socketLock)
+ {
+ //
+ if (socket != null)
+ {
+ //
+ socket.Close();
+ socket.Dispose();
+ socket = null;
+ }
+ }
+ }
+ catch { }
+
+ //
+ GC.SuppressFinalize(this);
+ }
+
+ //
+ #region Private ...
+ ///
+ /// Initialize Socket ...
+ ///
+ 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
+ }
+}
\ No newline at end of file