using System; using RabbitMQ.Client; using RabbitMQ.Client.Events; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Extensions; using xExceptions.Constants; namespace xEventService.Models { public class XChannel { // public IModel Channel { get; protected set; } public string QueueName { get; protected set; } public string ExchangeName { get; protected set; } public IConnection Connection { get; protected set; } // #region Constructors ... public XChannel ( string exchangeName = null, IConnection connection = null, IModel channel = null, EventingBasicConsumer consumer = null ) { // // Validate Connection ... if (connection.IsNull ()) { XException.InvalidArgs.Throw (); } Connection = connection; // // Validate Channel ... if (channel.IsNull ()) { channel = connection.CreateModel (); } Channel = channel; // // Check Queue Name and it's Declaration ... if (!exchangeName.IsNullOrEmpty ()) { // ExchangeName = exchangeName; } } public XChannel ( XEventServiceConfiguration configuration, bool dispatchConsumersAsync = true, string exchangeName = null ) { // // Generate Factory ... var factory = configuration.GetConnectionFactory (); if (dispatchConsumersAsync) { factory.DispatchConsumersAsync = true; } // // Create Cahnnel ... Connection = factory.CreateConnection (); Channel = Connection.CreateModel (); // // Check Exchange Declaration if provided ... if (!exchangeName.IsNullOrEmpty ()) { // ExchangeName = exchangeName; var exchangeArgs = configuration.GetExchangeArgs (ExchangeName); // // Declaring Exchange ... Channel.ExchangeDeclare ( exchange: exchangeName, type: exchangeArgs.Type, durable: exchangeArgs.Durable, autoDelete: exchangeArgs.AutoDelete ); // // Create a Queue ... var queueName = $"{exchangeName}[{Guid.NewGuid().ToString().GetDigits()}]"; var queueArgs = configuration.GetQueueArgs (queueName); QueueName = queueName; Channel.QueueDeclare ( queue: queueName, durable: queueArgs.Durable, exclusive: queueArgs.Exclusive, autoDelete: queueArgs.AutoDelete ); // // Bind Queue to Exchange ... Channel.QueueBind ( queue: queueName, exchange: exchangeName, routingKey: "" ); } } #endregion // #region Actions ... /// /// Close Channel ... /// public void CloseChannel () { // // Check Channel Exists ... if (!this.Channel.IsNull ()) { // // Close Channel if it's Open ... if (this.Channel.IsOpen) { this.Channel.Close (); } // this.Channel.Dispose (); this.Channel = null; } } /// /// Close Connection ... /// public void CloseConnection () { // // Check Connection Exists ... if (!this.Connection.IsNull ()) { // // Close Connection if it's Open ... if (this.Connection.IsOpen) { this.Connection.Close (); } // this.Connection.Dispose (); this.Connection = null; } } /// /// Publish a Message through Channel on Declare Queue ... /// /// /// /// public bool Publish (T message) { // // Validate Args ... if (this.Connection.IsNull () || !this.Connection.IsOpen || this.Channel.IsNull () || !this.Channel.IsOpen) { // Console.WriteLine ($"XEventService Exception: Publish Failed ..."); // XException.InvalidArgs.Throw (); } // // Check Publish Method ... // // Check Exchange Name ... if (this.ExchangeName.IsNullOrEmpty ()) { // Console.WriteLine ($"XEventService Exception: Publish Failed, Invalid Exchange Name {ExchangeName} ..."); // XException.InvalidArgs.Throw (); } // try { // // Publish Message through Channel to Declared Queue ... var body = message.ToBody (); this.Channel.BasicPublish ( exchange: ExchangeName, routingKey: "", basicProperties : null, body : body ); // return true; } catch (Exception ex) { // Console.WriteLine ($"XEventService Exception: Publish Failed, {ex.Message} ..."); // return false; } } public void Consume (EventingBasicConsumer consumer) { Channel.BasicConsume ( queue: QueueName, autoAck: true, consumer: consumer ); } public void ConsumeAsync (AsyncEventingBasicConsumer consumer) { Channel.BasicConsume ( queue: QueueName, autoAck: true, consumer: consumer ); } /// /// Dispose XChannel ... /// public void Dispose () { // this.CloseChannel (); this.CloseConnection (); } #endregion } }