diff --git a/Configuration/XEventServiceConfiguration.cs b/Configuration/XEventServiceConfiguration.cs
index 8e5df4c..f291fb0 100644
--- a/Configuration/XEventServiceConfiguration.cs
+++ b/Configuration/XEventServiceConfiguration.cs
@@ -1,5 +1,3 @@
-using System.Collections.Generic;
-using RabbitMQ.Client;
using xEventService.Constants;
using xEventService.Models;
diff --git a/Constants/XEventBroker.cs b/Constants/XEventBroker.cs
index 0d39d85..951a047 100644
--- a/Constants/XEventBroker.cs
+++ b/Constants/XEventBroker.cs
@@ -7,13 +7,16 @@ namespace xEventService.Constants
///
public enum XEventBroker
{
- [StringValue("None")]
+ [StringValue(XEventServiceConstants.None)]
None,
- [StringValue("XKafka")]
+ [StringValue(XEventServiceConstants.XKafka)]
XKafka,
- [StringValue("XRabbitMQ")]
- XRabbitMQ
+ [StringValue(XEventServiceConstants.XZeroMQ)]
+ XZeroMQ,
+
+ [StringValue(XEventServiceConstants.XRabbitMQ)]
+ XRabbitMQ,
}
}
\ No newline at end of file
diff --git a/Constants/XEventServiceConstants.cs b/Constants/XEventServiceConstants.cs
index e46022c..3843330 100644
--- a/Constants/XEventServiceConstants.cs
+++ b/Constants/XEventServiceConstants.cs
@@ -8,10 +8,18 @@ namespace xEventService.Constants
//
// Defaults ...
- public const string XDefaultTopic = "XSaherelm";
- public const string XDefaultConsumerGroup = "XSaherElmGroup";
+ public const string XDefaultTopic = "XSaherelm";
+ public const string XDefaultConsumerGroup = "XSaherElmGroup";
//
- public const string XWildcard = "*";
+ public const string None = "None";
+ public const string XKafka = "XKafka";
+ public const string PubSub = "PubSub";
+ public const string XZeroMQ = "XZeroMQ";
+ public const string PushPull = "PushPull";
+ public const string XRabbitMQ = "XRabbitMQ";
+
+ //
+ public const string XWildcard = "*";
}
}
\ No newline at end of file
diff --git a/Constants/XZeroMQPattern.cs b/Constants/XZeroMQPattern.cs
new file mode 100644
index 0000000..8b1b5ec
--- /dev/null
+++ b/Constants/XZeroMQPattern.cs
@@ -0,0 +1,16 @@
+using xExceptions.Attributes;
+
+namespace xEventService.Constants
+{
+ public enum XZeroMQPattern
+ {
+ [StringValue(XEventServiceConstants.None)]
+ None,
+
+ [StringValue(XEventServiceConstants.PubSub)]
+ PubSub,
+
+ [StringValue(XEventServiceConstants.PushPull)]
+ PushPull
+ }
+}
\ No newline at end of file
diff --git a/Extensions/XEventServiceExtensions.cs b/Extensions/XEventServiceExtensions.cs
index f6c1c29..22d6746 100644
--- a/Extensions/XEventServiceExtensions.cs
+++ b/Extensions/XEventServiceExtensions.cs
@@ -71,7 +71,7 @@ namespace xEventService.Extensions
var result =
!source.IsNullOrDefault() &&
!source.Exchange.IsNullOrEmpty() &&
- !source.ExchangeType.IsValid();
+ source.ExchangeType.IsValid();
//
return result;
@@ -110,7 +110,37 @@ namespace xEventService.Extensions
//
return result;
}
-
+
+ ///
+ /// Validate ZeroMQ Pattern ...
+ ///
+ ///
+ ///
+ public static bool IsValid(this XZeroMQPattern source)
+ {
+ //
+ var result = source != XZeroMQPattern.None;
+
+ //
+ return result;
+ }
+
+ ///
+ /// Validate ZeroMQ Options ...
+ ///
+ ///
+ ///
+ public static bool IsValid(this XZeroMQOptions source)
+ {
+ //
+ var result =
+ !source.IsNullOrDefault() &&
+ source.Pattern.IsValid();
+
+ //
+ return result;
+ }
+
///
/// Validate an Array is Wildcard passed or not ...
///
@@ -122,9 +152,9 @@ namespace xEventService.Extensions
var result =
source is not null &&
source.Length == 1 &&
- source.Any(s =>
+ source.Any(s =>
s == XEventServiceConstants.XWildcard);
-
+
//
return result;
}
@@ -145,6 +175,101 @@ namespace xEventService.Extensions
return result;
}
+ ///
+ /// Validate Sender ...
+ ///
+ ///
+ ///
+ ///
+ public static bool IsValidSender(
+ this XEventMessage source,
+ string[] allowedSenders
+ )
+ {
+ //
+ // Validate Source and Collection ...
+ var result =
+ source.IsValid() &&
+ allowedSenders
+ .IsAllowedCollection();
+ if (!result)
+ {
+ return result;
+ }
+
+ //
+ // Validate Sender ...
+ result =
+ allowedSenders.IsAllowedWildcard() ||
+ allowedSenders.Any(x =>
+ x.Equals(
+ source.Sender,
+ System.StringComparison.InvariantCultureIgnoreCase));
+
+ //
+ return result;
+ }
+
+ ///
+ /// Validate Action ...
+ ///
+ ///
+ ///
+ ///
+ public static bool IsValidAction(
+ this XEventMessage source,
+ string[] allowedActions
+ )
+ {
+ //
+ // Validate Source and Collection ...
+ var result =
+ source.IsValid() &&
+ allowedActions
+ .IsAllowedCollection();
+ if (!result)
+ {
+ return result;
+ }
+
+ //
+ // Validate Action ...
+ result =
+ allowedActions.IsAllowedWildcard() ||
+ allowedActions.Any(x =>
+ x.Equals(
+ source.Action,
+ System.StringComparison.InvariantCultureIgnoreCase));
+
+ //
+ return result;
+ }
+
+ ///
+ /// Completely Validate a Message, based on providing:
+ /// - Allowed Actions
+ /// - Allowed Senders
+ ///
+ ///
+ ///
+ ///
+ ///
+ public static bool IsValid(
+ this XEventMessage source,
+ string[] allowedSenders,
+ string[] allowedActions
+ )
+ {
+ //
+ var result =
+ source.IsValid() &&
+ source.IsValidAction(allowedActions) &&
+ source.IsValidSender(allowedSenders);
+
+ //
+ return result;
+ }
+
///
/// Parse Url and Extract Host address ...
///
diff --git a/Interfaces/IXKafkaProducerService.cs b/Interfaces/IXKafkaProducerService.cs
index e6f69b2..950d4f0 100644
--- a/Interfaces/IXKafkaProducerService.cs
+++ b/Interfaces/IXKafkaProducerService.cs
@@ -1,38 +1,8 @@
-using System;
-using System.Threading;
-using System.Threading.Tasks;
-using xEventService.Constants;
-using xEventService.Models;
-
namespace xEventService.Interfaces
{
///
/// a Service for Produce Kafka Messages ...
///
- public interface IXKafkaProducerService : IDisposable
- {
- ///
- /// Sending Message ...
- ///
- /// an Instance of to Provides Kafka Messaging requirement ...
- ///
- ///
- Task SendMessageAsync(
- XEventRequest request,
- CancellationToken cancellationToken = default
- );
-
- ///
- /// Sending a Message ...
- ///
- /// an Instance of to Provides Kafka Messaging requirement ...
- /// Topic for Sending Message ...
- ///
- ///
- Task SendMessageAsync(
- XEventMessage message,
- string topic = XEventServiceConstants.XDefaultTopic,
- CancellationToken cancellationToken = default
- );
- }
+ public interface IXKafkaProducerService : IXProducerServiceBase
+ { }
}
\ No newline at end of file
diff --git a/Interfaces/IXProducerServiceBase.cs b/Interfaces/IXProducerServiceBase.cs
new file mode 100644
index 0000000..6183bb7
--- /dev/null
+++ b/Interfaces/IXProducerServiceBase.cs
@@ -0,0 +1,35 @@
+using System;
+using System.Threading;
+using System.Threading.Tasks;
+using xEventService.Constants;
+using xEventService.Models;
+
+namespace xEventService.Interfaces
+{
+ public interface IXProducerServiceBase : IDisposable
+ {
+ ///
+ /// Sending Message ...
+ ///
+ /// an Instance of to Provides Kafka Messaging requirement ...
+ ///
+ ///
+ Task SendMessageAsync(
+ XEventRequest request,
+ CancellationToken cancellationToken = default
+ );
+
+ ///
+ /// Sending a Message ...
+ ///
+ /// an Instance of to Provides Kafka Messaging requirement ...
+ /// Topic for Sending Message ...
+ ///
+ ///
+ Task SendMessageAsync(
+ XEventMessage message,
+ string topic = XEventServiceConstants.XDefaultTopic,
+ CancellationToken cancellationToken = default
+ );
+ }
+}
\ No newline at end of file
diff --git a/Interfaces/IXRabbitMQProducerService.cs b/Interfaces/IXRabbitMQProducerService.cs
index 2bd0806..46bd0f4 100644
--- a/Interfaces/IXRabbitMQProducerService.cs
+++ b/Interfaces/IXRabbitMQProducerService.cs
@@ -8,22 +8,8 @@ namespace xEventService.Interfaces
///
/// RabbitMQ Producer Service Implementation ...
///
- public interface IXRabbitMQProducerService : IDisposable
+ public interface IXRabbitMQProducerService : IXProducerServiceBase
{
- ///
- /// Sending Message with Full Request Descriptor ...
- ///
- ///
- /// an Instance of .
- /// Topic is used as Exchange Name ...
- ///
- ///
- ///
- Task SendMessageAsync(
- XEventRequest request,
- CancellationToken cancellationToken = default
- );
-
///
/// Sending a Message to Default Exchange ...
///
diff --git a/Interfaces/IXZeroMQProducerService.cs b/Interfaces/IXZeroMQProducerService.cs
new file mode 100644
index 0000000..894d934
--- /dev/null
+++ b/Interfaces/IXZeroMQProducerService.cs
@@ -0,0 +1,7 @@
+namespace xEventService.Interfaces
+{
+ public interface IXZeroMQProducerService : IXProducerServiceBase
+ {
+
+ }
+}
\ No newline at end of file
diff --git a/Models/XRabbitMQOptions.cs b/Models/XRabbitMQOptions.cs
index 1b4584d..bbccf27 100644
--- a/Models/XRabbitMQOptions.cs
+++ b/Models/XRabbitMQOptions.cs
@@ -1,4 +1,3 @@
-using RabbitMQ.Client;
using xEventService.Constants;
namespace xEventService.Models
diff --git a/Models/XZeroMQOptions.cs b/Models/XZeroMQOptions.cs
new file mode 100644
index 0000000..e409c18
--- /dev/null
+++ b/Models/XZeroMQOptions.cs
@@ -0,0 +1,28 @@
+using xEventService.Constants;
+
+namespace xEventService.Models
+{
+ public class XZeroMQOptions
+ {
+ ///
+ /// Socket Pattern (PubSub, PushPull) ...
+ ///
+ public XZeroMQPattern Pattern { get; set; } = XZeroMQPattern.PubSub;
+
+ ///
+ /// Bind or Connect ...
+ /// True for Server (Bind), False for Client (Connect) ...
+ ///
+ public bool IsServer { get; set; } = false;
+
+ ///
+ /// High Water Mark for Sending ...
+ ///
+ public int SendHighWaterMark { get; set; } = 1000;
+
+ ///
+ /// High Water Mark for Receiving ...
+ ///
+ public int ReceiveHighWaterMark { get; set; } = 1000;
+ }
+}
\ No newline at end of file
diff --git a/Providers/XBaseEventConsumerBackgroundService.cs b/Providers/XBaseEventConsumerBackgroundService.cs
new file mode 100644
index 0000000..7dc94e5
--- /dev/null
+++ b/Providers/XBaseEventConsumerBackgroundService.cs
@@ -0,0 +1,68 @@
+using System;
+using Microsoft.Extensions.Logging;
+using xCommons.Extensions;
+using xCommons.Providers;
+using xEventService.Configuration;
+using xEventService.Constants;
+using xEventService.Extensions;
+using xExceptions.Constants;
+
+namespace xEventService.Providers
+{
+ ///
+ /// Base Event Consumer Service ...
+ ///
+ public abstract class XBaseEventConsumerBackgroundService : XResumableBackgroundService
+ {
+ ///
+ /// Consumed Topic Messages ...
+ ///
+ protected readonly string topic;
+
+ ///
+ /// Allowed Senders ...
+ ///
+ public readonly string[] allowedSenders;
+
+ ///
+ /// Allowed Actions ...
+ ///
+ public readonly string[] allowedActions;
+
+ ///
+ /// Configuration ...
+ ///
+ protected readonly XEventServiceConfiguration configuration;
+
+ protected XBaseEventConsumerBackgroundService(
+ string[] allowedActions,
+ string[] allowedSenders,
+ IServiceProvider serviceProvider,
+ XEventServiceConfiguration configuration,
+ ILogger logger,
+ string topic = XEventServiceConstants.XDefaultTopic
+ ) : base(
+ logger: logger,
+ serviceProvider: serviceProvider
+ )
+ {
+ //
+ // Validate Actions ...
+ if (!allowedActions
+ .IsAllowedCollection() ||
+ !allowedSenders
+ .IsAllowedCollection() ||
+ !configuration.IsValid()
+ )
+ {
+ XException.InvalidArgs.Throw();
+ }
+
+ //
+ this.topic = topic;
+ this.configuration = configuration;
+ this.allowedActions = allowedActions;
+ this.allowedSenders = allowedSenders;
+ }
+ }
+}
\ No newline at end of file
diff --git a/Providers/XKafkaConsumerServiceBase.cs b/Providers/XKafkaConsumerServiceBase.cs
index 5ed8520..a55d452 100644
--- a/Providers/XKafkaConsumerServiceBase.cs
+++ b/Providers/XKafkaConsumerServiceBase.cs
@@ -2,12 +2,10 @@ using System;
using System.Threading;
using System.Threading.Tasks;
using Confluent.Kafka;
-using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
-using xEventService.Extensions;
using xExceptions.Constants;
namespace xEventService.Providers
@@ -15,67 +13,35 @@ namespace xEventService.Providers
///
/// an Abstract Kafka Consumer Service ...
///
- public abstract class XKafkaConsumerServiceBase : BackgroundService
+ public abstract class XKafkaConsumerServiceBase : XBaseEventConsumerBackgroundService
{
- ///
- /// Consumed Topic Messages ...
- ///
- protected readonly string topic;
-
- ///
- /// Check is Paused ...
- ///
- private bool paused = false;
-
- ///
- /// Allowed Senders ...
- ///
- public readonly string[] allowedSenders;
-
- ///
- /// Allowed Actions ...
- ///
- public readonly string[] allowedActions;
-
///
/// Kafka Consumer Class ...
///
protected readonly IConsumer consumer;
- protected readonly XEventServiceConfiguration configuration;
-
- protected readonly ILogger logger;
-
protected XKafkaConsumerServiceBase(
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
)
{
//
- this.topic = topic;
- this.logger = logger;
- this.configuration = configuration;
-
- //
- // Validate Actions ...
- if (!allowedActions
- .IsAllowedCollection())
+ // Validating Configuration ...
+ if (configuration.Broker != XEventBroker.XKafka)
{
- XException.InvalidArgs.Throw();
+ XException.InvalidConfiguration.Throw();
}
- this.allowedActions = [.. allowedActions];
-
- //
- // Validate Senders ...
- if (!allowedSenders
- .IsAllowedCollection())
- {
- XException.InvalidArgs.Throw();
- }
- this.allowedSenders = [.. allowedSenders];
//
// Prepare Kafka Consumer Configuration ...
@@ -89,45 +55,6 @@ namespace xEventService.Providers
consumer = new ConsumerBuilder(config).Build();
}
-
- ///
- /// Pause Consuming ...
- ///
- ///
- public bool Pause()
- {
- //
- var result = !paused;
- if (!result)
- {
- return result;
- }
-
- //
- paused = true;
- result = paused;
- return result;
- }
-
- ///
- /// Resume Consuming ...
- ///
- ///
- public bool Resume()
- {
- //
- var result = paused;
- if (!result)
- {
- return result;
- }
-
- //
- paused = false;
- result = !paused;
- return result;
- }
-
///
/// Default Action Execution for Background Services ...
///
@@ -142,94 +69,40 @@ namespace xEventService.Providers
await Task.Yield();
using var cts = new CancellationTokenSource();
+ //
// Subscribe to Topic ...
consumer.Subscribe(topic);
+ //
// Consuming ...
while (!cancellationToken.IsCancellationRequested)
{
try
{
+ //
// Consuming Topic ...
- await consumer.Consume(async (consumeResult, cancellationToken) =>
+ var consumeResult = consumer.Consume(cancellationToken);
+ if (!consumeResult.IsNullOrDefault() &&
+ !consumeResult.Message.IsNullOrDefault()
+ )
{
- // Check Pausing ...
- if (!paused)
- {
- // Logging Message Recieved ...
- if (enableLogging)
- {
- // Encrypt Message ...
- var expirationOffset = DateTimeOffset.UtcNow
- .AddMinutes(encryptionLifetimeInMinute);
- var encryptedMessage = cryptographyService.Encrypt(
- text: consumeResult?.Message.Value ?? string.Empty,
- expirationOffset: expirationOffset
- );
- logger
- .LogInformation(
- "Consume Recieved Message from Kafka Topic: {Topic}, Offset: {Offset}, Partition: {Partition}, Message: {Message}",
- topic,
- consumeResult?.Offset,
- consumeResult?.Partition,
- encryptedMessage
- );
- }
-
- // Produce Message Model ...
- var model = consumeResult
- .FillConsumeResult()
- ?? throw RpkKafkaExceptionHelper
- .GetInvalidKafkaConsumeResultException();
-
- // Check Sender and Actions Owning ...
- var isOwned = model.IsOwned(
- allowedActions,
- allowedSenders
- );
- if (isOwned)
- {
- // Notify Extended Classes for Consume Message ...
- await ConsumeAsync(
- model,
- cancellationToken
- );
- }
-
- // Commit Offset after Consuming ...
- consumer?.Commit(consumeResult);
- }
- },
- cts.Token);
- }
- catch (ConsumeException ex)
- {
- // Log Exception ...
- if (enableLogging)
- {
- logger.LogError(
- ex,
- "Kafka Message Consume Failed: {Error}",
- ex.Error.Reason
+ //
+ await DoWorkAsync(
+ payload: consumeResult.Message,
+ cancellationToken: cancellationToken
);
}
+ }
+ catch (ConsumeException)
+ {
await Task.Delay(1000, cancellationToken);
}
- catch (Exception ex)
+ catch (Exception)
{
- // Log Exception ...
- if (enableLogging)
- {
- logger.LogError(
- ex,
- "Exception: {Exc}",
- ex.Message
- );
- }
await Task.Delay(1000, cancellationToken);
}
}
- }
+ }
///
/// Dispose Implementation ...
diff --git a/Providers/XKafkaProducerService.cs b/Providers/XKafkaProducerService.cs
index 1bf4da0..0c693fb 100644
--- a/Providers/XKafkaProducerService.cs
+++ b/Providers/XKafkaProducerService.cs
@@ -115,7 +115,7 @@ namespace xEventService.Providers
);
//
- result = true;
+ result = !response.IsNullOrDefault();
}
catch
{
diff --git a/Providers/XRabbitMQConsumerServiceBase.cs b/Providers/XRabbitMQConsumerServiceBase.cs
index 1fa5b72..37ad228 100644
--- a/Providers/XRabbitMQConsumerServiceBase.cs
+++ b/Providers/XRabbitMQConsumerServiceBase.cs
@@ -1,13 +1,15 @@
using System;
-using System.Collections.Generic;
-using System.Linq;
+using System.Text;
+using System.Threading;
using System.Threading.Tasks;
-using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
+using RabbitMQ.Client.Events;
using xCommons.Extensions;
using xEventService.Configuration;
+using xEventService.Constants;
using xEventService.Extensions;
+using xEventService.Models;
using xExceptions.Constants;
namespace xEventService.Providers
@@ -15,36 +17,233 @@ namespace xEventService.Providers
///
/// an Abstract Kafka Consumer Service ...
///
- public class XRabbitMQConsumerServiceBase : BackgroundService
+ public abstract class XRabbitMQConsumerServiceBase : XBaseEventConsumerBackgroundService
{
- protected readonly IModel channel;
- protected readonly IConnection connection;
- protected readonly XEventServiceConfiguration configuration;
- protected readonly ILogger logger;
+ protected IModel channel;
+ protected IConnection connection;
+ protected AsyncEventingBasicConsumer consumer;
- public XRabbitMQConsumerServiceBase(
+ protected XRabbitMQConsumerServiceBase(
+ string[] allowedActions,
+ string[] allowedSenders,
+ IServiceProvider serviceProvider,
XEventServiceConfiguration configuration,
- ILogger logger
+ ILogger logger,
+ string topic = XEventServiceConstants.XDefaultTopic
+ ) : base(
+ topic: topic,
+ logger: logger,
+ configuration: configuration,
+ allowedActions: allowedActions,
+ allowedSenders: allowedSenders,
+ serviceProvider: serviceProvider
)
{
//
- this.logger = logger;
- this.configuration = configuration;
-
- //
- // Validate ...
- if (!configuration.IsValid() ||
- !configuration.RabbitMQ.IsValid() ||
- configuration.Broker != Constants.XEventBroker.XRabbitMQ
+ // Validating Configuration ...
+ if (!configuration.RabbitMQ.IsValid() ||
+ configuration.Broker != XEventBroker.XRabbitMQ
)
{
XException.InvalidConfiguration.Throw();
}
+
+ //
+ SetupConnection();
+ DeclareTopology();
+ consumer = new AsyncEventingBasicConsumer(channel);
}
-
- public abstract Task HandleMessageAsync(
- XEventMessage message,
+
+ ///
+ /// Default Action Execution for Background Services ...
+ ///
+ ///
+ ///
+ protected override async Task ExecuteAsync(
CancellationToken cancellationToken = default
- );
+ )
+ {
+ //
+ // For Fix Blocking Synchronous ...
+ await Task.Yield();
+ using var cts = new CancellationTokenSource();
+
+ //
+ // Subscribe to Topic ...
+ consumer.Received += async (model, ea) =>
+ {
+ //
+ try
+ {
+ //
+ var body = ea.Body.ToArray();
+ var json = Encoding.UTF8.GetString(body);
+ var message = json.FromJSON();
+ if (!message.IsNullOrDefault())
+ {
+ //
+ await DoWorkAsync(
+ payload: message,
+ cancellationToken: cancellationToken
+ );
+ }
+
+ //
+ if (configuration.RabbitMQ.ManualAck)
+ {
+ //
+ channel.BasicAck(
+ multiple: false,
+ deliveryTag: ea.DeliveryTag
+ );
+ }
+ }
+ catch
+ {
+ //
+ if (configuration.RabbitMQ.ManualAck)
+ {
+ //
+ channel.BasicNack(
+ requeue: true,
+ multiple: false,
+ deliveryTag: ea.DeliveryTag
+ );
+ }
+ }
+ };
+
+ //
+ consumer.Shutdown += (model, ea) =>
+ {
+ return Task.CompletedTask;
+ };
+
+ //
+ channel.BasicQos(
+ prefetchSize: 0,
+ global: false,
+ prefetchCount: configuration.RabbitMQ.PrefetchCount
+ );
+
+ //
+ channel.BasicConsume(
+ queue: GetQueueName(),
+ consumer: consumer,
+ autoAck: !configuration.RabbitMQ.ManualAck
+ );
+
+ //
+ while (!cancellationToken.IsCancellationRequested)
+ {
+ await Task.Delay(1000, cancellationToken);
+ }
+ }
+
+ ///
+ /// Overriding Stop Action ...
+ ///
+ ///
+ ///
+ public override async Task StopAsync(CancellationToken cancellationToken)
+ {
+ //
+ channel.Close();
+ channel.Dispose();
+ connection.Close();
+ connection.Dispose();
+
+ //
+ await base.StopAsync(cancellationToken);
+ }
+
+ ///
+ /// Dispose Implementation ...
+ ///
+ public override void Dispose()
+ {
+ //
+ channel.Dispose();
+ connection.Dispose();
+ GC.SuppressFinalize(this);
+ base.Dispose();
+ }
+
+ //
+ #region Private ...
+ private string GetQueueName()
+ {
+ //
+ var result = string.Empty;
+ if (!configuration.RabbitMQ.Queue.IsNullOrEmpty())
+ {
+ result = configuration.RabbitMQ.Queue;
+ }
+
+ //
+ if (result.IsNullOrEmpty())
+ {
+ //
+ var digits = Guid.NewGuid().ToString().GetDigits();
+ result = $"{configuration.RabbitMQ.Exchange}[{digits}]";
+ }
+
+ //
+ return result;
+ }
+
+ private void SetupConnection()
+ {
+ var factory = new ConnectionFactory
+ {
+ DispatchConsumersAsync = true,
+ Port = configuration.GetRabbitMQPort(),
+ HostName = configuration.GetRabbitMQHost(),
+ UserName = configuration.RabbitMQ.Username,
+ Password = configuration.RabbitMQ.Password,
+ NetworkRecoveryInterval = TimeSpan.FromSeconds(
+ configuration.RabbitMQ.NetworkRecoveryIntervalSeconds
+ ),
+ VirtualHost = configuration.RabbitMQ.VirtualHost,
+ AutomaticRecoveryEnabled = configuration.RabbitMQ.AutomaticRecoveryEnabled,
+ };
+
+ //
+ connection = factory.CreateConnection();
+ channel = connection.CreateModel();
+ }
+
+ private void DeclareTopology()
+ {
+ //
+ channel.ExchangeDeclare(
+ arguments: null,
+ exchange: configuration.RabbitMQ.Exchange,
+ durable: configuration.RabbitMQ.ExchangeDurable,
+ autoDelete: configuration.RabbitMQ.ExchangeAutoDelete,
+ type: configuration.RabbitMQ.ExchangeType.GetStringValue()
+ );
+
+ //
+ var queueName = GetQueueName();
+
+ //
+ channel.QueueDeclare(
+ arguments: null,
+ queue: queueName,
+ durable: configuration.RabbitMQ.QueueDurable,
+ exclusive: configuration.RabbitMQ.QueueExclusive,
+ autoDelete: configuration.RabbitMQ.QueueAutoDelete
+ );
+
+ //
+ channel.QueueBind(
+ arguments: null,
+ queue: queueName,
+ exchange: configuration.RabbitMQ.Exchange,
+ routingKey: configuration.RabbitMQ.RoutingKey
+ );
+ }
+ #endregion
}
}
\ No newline at end of file
diff --git a/Providers/XRabbitMQProducerService.cs b/Providers/XRabbitMQProducerService.cs
index 74a7b80..3d27b63 100644
--- a/Providers/XRabbitMQProducerService.cs
+++ b/Providers/XRabbitMQProducerService.cs
@@ -231,7 +231,7 @@ namespace xEventService.Providers
///
/// Ensure Connection Established ...
///
- public async Task EnsureConnectionAsync()
+ private async Task EnsureConnectionAsync()
{
//
try
@@ -289,9 +289,6 @@ namespace xEventService.Providers
type: configuration.RabbitMQ.ExchangeType.GetStringValue()
);
}
-
- //
- await Task.CompletedTask;
}
#endregion
}
diff --git a/xEventService.csproj b/xEventService.csproj
index ebc3ecc..73a4f3b 100644
--- a/xEventService.csproj
+++ b/xEventService.csproj
@@ -20,8 +20,7 @@
-
+
@@ -35,6 +34,7 @@
+