### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Configuration\XEventServiceConfiguration.cs using System.Collections.Generic; using RabbitMQ.Client; using xEventService.Constants; using xEventService.Models; namespace xEventService.Configuration { /// /// a configuration class for describing how to connect rabbit mq server /// and how channels must be declared ... /// public class XEventServiceConfiguration { /// /// Message Broker Address ... /// public string Url { get; set; } /// /// Max Retries to Sending Message ... /// public int MaxRetries { get; set; } = 3; /// /// Allowed Message Broker ... /// public XEventBroker Broker { get; set; } = XEventBroker.None; /// /// RabbitMQ Specific Options ... /// Only used when Broker is XRabbitMQ ... /// public XRabbitMQOptions RabbitMQ { get; set; } = new XRabbitMQOptions(); } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\ConfigurationNodeNames.cs namespace xEventService.Constants { public partial struct ConfigurationNodeNames { public const string EVENT_SERVICE_NODE_NAME = "EventServiceConfiguration"; } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\XEventBroker.cs using xExceptions.Attributes; namespace xEventService.Constants { /// /// Available Brokers ... /// public enum XEventBroker { [StringValue("None")] None, [StringValue("XKafka")] XKafka, [StringValue("XRabbitMQ")] XRabbitMQ } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\XEventServiceConstants.cs namespace xEventService.Constants { public struct XEventServiceConstants { // // Service Extensions Log Tag ... public const string XEventServiceDILogTag = "XEventService"; // // Defaults ... public const string XDefaultTopic = "XSaherelm"; public const string XDefaultConsumerGroup = "XSaherElmGroup"; // public const string XWildcard = "*"; } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\XRabbitMQExchangeType.cs using xExceptions.Attributes; namespace xEventService.Constants { /// /// Available RabbitMQ Exchange Types ... /// public struct XRabbitMQExchangeTypes { public const string None = "none"; public const string Direct = "direct"; public const string Fanout = "fanout"; public const string Topic = "topic"; public const string Headers = "headers"; } /// /// RabbitMQ Exchange Type Enum ... /// public enum XRabbitMQExchangeType { [StringValue(XRabbitMQExchangeTypes.None)] None, [StringValue(XRabbitMQExchangeTypes.Direct)] Direct, [StringValue(XRabbitMQExchangeTypes.Fanout)] Fanout, [StringValue(XRabbitMQExchangeTypes.Topic)] Topic, [StringValue(XRabbitMQExchangeTypes.Headers)] Headers, } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\DI\XDIHelperExtension.cs using System; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xEventService.Extensions; using xExceptions.Constants; namespace xEventService.DI { public static partial class XDIHelperExtension { /// /// Extract EventService Configurations ... /// /// /// public static XEventServiceConfiguration GetXEventServiceConfigurations( this IConfiguration source ) { // var configSection = source .GetSection(ConfigurationNodeNames.EVENT_SERVICE_NODE_NAME); var result = configSection.Get(); // return result; } /// /// Register Event Service Configurations ... /// /// /// public static void AddXEventServiceConfigurations( this IServiceCollection services, IConfiguration configuration ) { // var config = configuration.GetXEventServiceConfigurations(); services.AddXEventServiceConfigurations(config); } /// /// Register Event Service Configurations ... /// /// /// public static void AddXEventServiceConfigurations( this IServiceCollection services, XEventServiceConfiguration configuration ) { // // Validate ... if (!configuration.IsValid()) { // Log("Registration Failed, due Invalid Configuration ..."); XException.InvalidConfiguration.Throw(); } // // Register Configuration as Singleton ... services.AddSingleton(configuration); } // #region Private ... /// /// a LogTag for Service ... /// private static string XLogTag = XEventServiceConstants.XEventServiceDILogTag; /// /// print a log in Console ... /// /// private static void Log(string message) { Console.WriteLine($"{XLogTag} => {message}"); } #endregion } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Extensions\XEventServiceExtensions.cs using System.Linq; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xEventService.Models; using xExceptions.Constants; namespace xEventService.Extensions { /// /// Extensions for Event Service ... /// public static class XEventServiceExtensions { /// /// Validate Event Broker ... /// /// /// public static bool IsValid(this XEventBroker source) { // var result = !source.IsNull() && source != XEventBroker.None; // return result; } /// /// Validate Event Service Configuration ... /// /// /// public static bool IsValid(this XEventServiceConfiguration source) { // var result = !source.IsNullOrDefault() && source.Broker.IsValid() && !source.Url.IsNullOrEmpty() && source.Url.IsValidUrl(); // return result; } /// /// Validate RabbitMQ Exchange Type ... /// /// /// public static bool IsValid(this XRabbitMQExchangeType source) { // var result = source != XRabbitMQExchangeType.None; // return result; } /// /// Validate RabbitMQ Configuration ... /// /// /// public static bool IsValid(this XRabbitMQOptions source) { // var result = !source.IsNullOrDefault() && !source.Exchange.IsNullOrEmpty() && source.ExchangeType.IsValid(); // return result; } /// /// Validate an Event Message ... /// /// /// public static bool IsValid(this XEventMessage source) { // var result = !source.IsNullOrDefault() && !source.Action.IsNullOrEmpty() && !source.Sender.IsNullOrEmpty(); // return result; } /// /// Validate an Event Request ... /// /// /// public static bool IsValid(this XEventRequest source) { // var result = !source.IsNullOrDefault() && !source.Topic.IsNullOrEmpty() && source.Message.IsValid(); // return result; } /// /// Validate an Array is Wildcard passed or not ... /// /// /// public static bool IsAllowedWildcard(this string[] source) { // var result = source is not null && source.Length == 1 && source.Any(s => s == XEventServiceConstants.XWildcard); // return result; } /// /// Validate a Collection is Valid for Kafka Consumer ... /// /// /// public static bool IsAllowedCollection( this string[] source ) { // var result = source.HasChild(); // return result; } /// /// Parse Url and Extract Host address ... /// /// /// public static string GetRabbitMQHost(this XEventServiceConfiguration source) { // var result = string.Empty; // if (!source.IsValid()) { XException.InvalidConfiguration.Throw(); } // var parts = source.Url.Split(':'); if (parts.Length != 2) { XException.InvalidConfiguration.Throw(); } // result = parts[0]; if (result.IsNullOrEmpty() || !result.IsValidUrl() ) { XException.InvalidConfiguration.Throw(); } // return result; } /// /// Parse Url and Extract Port Number ... /// /// /// public static int GetRabbitMQPort(this XEventServiceConfiguration source) { // int result = -1; // if (!source.IsValid()) { XException.InvalidConfiguration.Throw(); } // var parts = source.Url.Split(':'); if (parts.Length != 2) { XException.InvalidConfiguration.Throw(); } // // Try Parse Port as Integer ... int.TryParse(parts[1], out result); if (result <= 0) { XException.InvalidConfiguration.Throw(); } // return result; } } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXKafkaProducerService.cs 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 ); } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXRabbitMQProducerService.cs using System; using System.Threading; using System.Threading.Tasks; using xEventService.Models; namespace xEventService.Interfaces { /// /// RabbitMQ Producer Service Implementation ... /// public interface IXRabbitMQProducerService : IDisposable { /// /// 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 ... /// /// /// an Instance of ... /// /// /// Target Exchange Name. if empty, uses default from configuration ... /// /// /// Routing Key for message routing ... /// /// /// Task SendMessageAsync( XEventMessage message, string exchange = null, string routingKey = null, CancellationToken cancellationToken = default ); } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Models\XEventMessage.cs using System; using System.Collections.Generic; using Newtonsoft.Json; namespace xEventService.Models { /// /// a Message Descriptor for Events ... /// public class XEventMessage { /// /// Specified Action for Fire ... /// [JsonRequired] public string Action { get; set; } /// /// Specified Action Sender ... /// [JsonRequired] public string Sender { get; set; } /// /// Sending Time ... /// [JsonRequired] public DateTime Offset { get; set; } /// /// Specified Message ... /// public string Message { get; set; } /// /// Metadata for Message ... /// public IDictionary Payload { get; set; } = new Dictionary(); } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Models\XEventRequest.cs using xEventService.Constants; namespace xEventService.Models { /// /// a Request for Sending an Event ... /// public class XEventRequest { /// /// Specified Which Kafka Topic ... /// public string Topic { get; set; } = XEventServiceConstants.XDefaultTopic; /// /// Specified Kafka Meesage to Send ... /// public XEventMessage Message { get; set; } } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Models\XRabbitMQOptions.cs using RabbitMQ.Client; using xEventService.Constants; namespace xEventService.Models { /// /// RabbitMQ-specific connection and messaging options ... /// public class XRabbitMQOptions { /// /// Connection Username ... /// public string Username { get; set; } = "guest"; /// /// Connection Password ... /// public string Password { get; set; } = "guest"; /// /// Virtual Host ... /// public string VirtualHost { get; set; } = "/"; /// /// Default Exchange Name ... /// public string Exchange { get; set; } = "x.saherelm.events"; /// /// Default Exchange Type (fanout, direct, topic, headers) ... /// public XRabbitMQExchangeType ExchangeType { get; set; } = XRabbitMQExchangeType.Fanout; /// /// Default Queue Name ... /// if empty, a dynamic queue will be generated ... /// public string Queue { get; set; } = string.Empty; /// /// Default Routing Key ... /// public string RoutingKey { get; set; } = string.Empty; /// /// Exchange Durable Flag ... /// public bool ExchangeDurable { get; set; } = true; /// /// Exchange AutoDelete Flag ... /// public bool ExchangeAutoDelete { get; set; } = false; /// /// Queue Durable Flag ... /// public bool QueueDurable { get; set; } = true; /// /// Queue Exclusive Flag ... /// public bool QueueExclusive { get; set; } = false; /// /// Queue AutoDelete Flag ... /// public bool QueueAutoDelete { get; set; } = false; /// /// Enable Automatic Recovery on Connection Loss ... /// public bool AutomaticRecoveryEnabled { get; set; } = true; /// /// Network Recovery Interval in Seconds ... /// public int NetworkRecoveryIntervalSeconds { get; set; } = 5; /// /// Prefetch Count for Consumer ... /// public ushort PrefetchCount { get; set; } = 1; /// /// Enable Manual Acknowledgement ... /// public bool ManualAck { get; set; } = true; } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XBaseEventConsumerBackgroundService.cs 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; } } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XKafkaConsumerServiceBase.cs using System; using System.Threading; using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.Logging; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xExceptions.Constants; namespace xEventService.Providers { /// /// an Abstract Kafka Consumer Service ... /// public abstract class XKafkaConsumerServiceBase : XBaseEventConsumerBackgroundService { /// /// Kafka Consumer Class ... /// protected readonly IConsumer consumer; 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 ) { // // Validating Configuration ... if (configuration.Broker != XEventBroker.XKafka) { XException.InvalidConfiguration.Throw(); } // // Prepare Kafka Consumer Configuration ... var config = new ConsumerConfig { EnableAutoCommit = false, BootstrapServers = configuration.Url, AutoOffsetReset = AutoOffsetReset.Earliest, GroupId = XEventServiceConstants.XDefaultConsumerGroup, }; consumer = new ConsumerBuilder(config).Build(); } /// /// 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.Subscribe(topic); // // Consuming ... while (!cancellationToken.IsCancellationRequested) { try { // // Consuming Topic ... var consumeResult = consumer.Consume(cancellationToken); if (consumeResult.IsNullOrDefault() || !consumeResult.Message.IsNullOrDefault() ) { // await DoWorkAsync( payload: consumeResult.Message, cancellationToken: cancellationToken ); } } catch (ConsumeException) { await Task.Delay(1000, cancellationToken); } catch (Exception) { await Task.Delay(1000, cancellationToken); } } } /// /// Dispose Implementation ... /// public override void Dispose() { // consumer.Close(); consumer.Dispose(); GC.SuppressFinalize(this); base.Dispose(); } } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XKafkaProducerService.cs using xCommons.Extensions; using xExceptions.Constants; using xEventService.Extensions; using xEventService.Interfaces; using xEventService.Configuration; using xEventService.Models; using Confluent.Kafka; using Microsoft.Extensions.Logging; using System.Threading.Tasks; using System.Threading; using xEventService.Constants; using System; namespace xEventService.Providers { public class XKafkaProducerService : IXKafkaProducerService { private readonly IProducer producer; private readonly ILogger logger; private readonly XEventServiceConfiguration configuration; public XKafkaProducerService( ILogger logger, XEventServiceConfiguration configuration ) { // this.logger = logger; this.configuration = configuration; // // Validate Configurations ... if (!configuration.IsValid() || configuration.Broker != XEventBroker.XKafka ) { XException.InvalidConfiguration.Throw(); } // // Prepare Kafka Producer Configuration ... var config = new ProducerConfig { Acks = Acks.All, BootstrapServers = configuration.Url, MessageSendMaxRetries = configuration.MaxRetries, }; producer = new ProducerBuilder(config).Build(); } /// /// Sending Message ... /// /// an Instance of to Provides Kafka Messaging requirement ... /// /// public async Task SendMessageAsync( XEventRequest request, CancellationToken cancellationToken = default ) { // // Validate ... 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 ... /// /// an Instance of to Provides Kafka Messaging requirement ... /// Topic for Sending Message ... /// /// public async Task SendMessageAsync( XEventMessage message, string topic = XEventServiceConstants.XDefaultTopic, CancellationToken cancellationToken = default ) { // // Validate ... var result = message.IsValid() && !topic.IsNullOrEmpty(); if (!result) { XException.InvalidArgs.Throw(); } // try { // var response = await producer.ProduceAsync( topic: topic, message: new Message { Value = message .ToJSON(camelCase: true) }, cancellationToken: cancellationToken ); // result = !response.IsNullOrDefault(); } catch { result = false; } // return result; } /// /// Dispose Required Objects ... /// public void Dispose() { // producer.Dispose(); GC.SuppressFinalize(this); } } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XRabbitMQConsumerServiceBase.cs using System; using System.Text; using System.Threading; using System.Threading.Tasks; 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 { /// /// an Abstract Kafka Consumer Service ... /// public abstract class XRabbitMQConsumerServiceBase : XBaseEventConsumerBackgroundService { protected IModel channel; protected IConnection connection; protected AsyncEventingBasicConsumer consumer; protected XRabbitMQConsumerServiceBase( 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.RabbitMQ.IsValid() || configuration.Broker != XEventBroker.XRabbitMQ ) { XException.InvalidConfiguration.Throw(); } // SetupConnection(); DeclareTopology(); consumer = new AsyncEventingBasicConsumer(channel); } /// /// 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 } } ### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XRabbitMQProducerService.cs using System; using System.Text; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.Logging; using RabbitMQ.Client; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xEventService.Extensions; using xEventService.Interfaces; using xEventService.Models; using xExceptions.Constants; namespace xEventService.Providers { /// /// RabbitMQ Producer Service Implementation ... /// public class XRabbitMQProducerService : 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 ) { // this.logger = logger; this.configuration = configuration; // // Validate Configurations ... if (!configuration.IsValid() || !configuration.RabbitMQ.IsValid() || configuration.Broker != XEventBroker.XRabbitMQ ) { XException.InvalidConfiguration.Throw(); } // // Initialize Connection ... CreateConnection(); } /// /// Sending Message with Full Request Descriptor ... /// public async Task SendMessageAsync( XEventRequest request, CancellationToken cancellationToken = default ) { // // Validate ... var result = request.IsValid(); if (!result) { XException.InvalidArgs.Throw(); } // result = await SendMessageAsync( message: request.Message, exchange: request.Topic, cancellationToken: cancellationToken, routingKey: configuration.RabbitMQ.RoutingKey ); // return result; } /// /// Sending a Message to Exchange ... /// public async Task SendMessageAsync( XEventMessage message, string exchange = null, string routingKey = null, CancellationToken cancellationToken = default ) { // // Validate ... var result = message.IsValid(); if (!result) { XException.InvalidArgs.Throw(); } // // Normalize Exchange ... if (exchange.IsNullOrEmpty()) { exchange = configuration.RabbitMQ.Exchange; } // // Normalize Routing Key ... if (routingKey.IsNull()) { routingKey = configuration.RabbitMQ.RoutingKey; } // try { // // Ensure Connection is Alive ... await EnsureChannelAsync(cancellationToken); // // Serialize Message ... var body = Encoding.UTF8.GetBytes( message.ToJSON(camelCase: true) ); // // Prepare Message Properties ... var properties = channel.CreateBasicProperties(); properties.Persistent = true; properties.ContentType = "application/json"; properties.Timestamp = new AmqpTimestamp( DateTimeOffset.UtcNow.ToUnixTimeSeconds() ); properties.MessageId = Guid.NewGuid().ToString(); properties.Type = message.Action; properties.AppId = message.Sender; // // Publish Message ... await Task.Run(() => { // channel.BasicPublish( exchange: exchange, routingKey: routingKey, mandatory: false, basicProperties: properties, body: body ); }, cancellationToken: cancellationToken); // result = true; } catch { result = false; } // return result; } /// /// Dispose Object ... /// public void Dispose() { // try { // channel.Close(); channel.Dispose(); channel = null; // connection.Close(); connection.Dispose(); connection = null; } catch { } // GC.SuppressFinalize(this); } // #region Private ... /// /// Create Connection to RabbitMQ Server ... /// private void CreateConnection() { // var factory = new ConnectionFactory { Port = configuration.GetRabbitMQPort(), HostName = configuration.GetRabbitMQHost(), UserName = configuration.RabbitMQ.Username, Password = configuration.RabbitMQ.Password, VirtualHost = configuration.RabbitMQ.VirtualHost, AutomaticRecoveryEnabled = configuration.RabbitMQ.AutomaticRecoveryEnabled, NetworkRecoveryInterval = TimeSpan.FromSeconds( configuration.RabbitMQ.NetworkRecoveryIntervalSeconds ), DispatchConsumersAsync = true }; // connection = factory.CreateConnection(); } /// /// Check Connection is Alive ... /// private bool IsConnected() { // var result = !connection.IsNullOrDefault() && connection.IsOpen && !channel.IsNull() && channel.IsOpen; // 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 ) { // if (IsConnected()) { return; } // lock (connectionLock) { // // Double Check after Lock ... if (IsConnected()) { return; } // // Recreate Connection if Dead ... if (connection is null || !connection.IsOpen) { CreateConnection(); } // // Create Channel ... channel = connection.CreateModel(); // // Declare Exchange ... channel.ExchangeDeclare( arguments: null, exchange: configuration.RabbitMQ.Exchange, durable: configuration.RabbitMQ.ExchangeDurable, autoDelete: configuration.RabbitMQ.ExchangeAutoDelete, type: configuration.RabbitMQ.ExchangeType.GetStringValue() ); } // await Task.CompletedTask; } #endregion } }