using System; using System.Text; using RabbitMQ.Client; using System.Threading; using xCommons.Extensions; using xEventService.Models; using xExceptions.Constants; using System.Threading.Tasks; using xEventService.Constants; using xEventService.Extensions; using xEventService.Interfaces; using xEventService.Configuration; using Microsoft.Extensions.Logging; namespace xEventService.Providers { /// /// RabbitMQ Producer Service Implementation ... /// public class XRabbitMQProducerService : XBaseEventProducerService, IXRabbitMQProducerService { private IModel channel; private IConnection connection; private readonly object connectionLock = new(); public XRabbitMQProducerService( ILogger logger, XEventServiceConfiguration configuration ) : base(configuration, logger) { // // 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 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, 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 ... EnsureChannel(); // // 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 (Exception ex) { // logger.LogError($"RabbitMQ Produce Error: {ex.Message} ..."); result = false; } // return result; } /// /// Dispose Object ... /// public void Dispose() { // try { // if (channel != null) { // channel.Close(); channel.Dispose(); channel = null; } // if (connection != 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 Channel is Open and Exchange/Queue Declared ... /// private void EnsureChannel() { // 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() ); } } #endregion } }