### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Configuration\XEventServiceConfiguration.cs 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; /// /// 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 ... /// public XRabbitMQOptions RabbitMQ { get; set; } = new XRabbitMQOptions(); } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Constants\ConfigurationNodeNames.cs namespace xEventService.Constants { public partial struct ConfigurationNodeNames { public const string EVENT_SERVICE_NODE_NAME = "EventServiceConfiguration"; } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Constants\XEventBroker.cs using xExceptions.Attributes; namespace xEventService.Constants { /// /// Available Brokers ... /// public enum XEventBroker { [StringValue(XEventServiceConstants.None)] None, [StringValue(XEventServiceConstants.XKafka)] XKafka, [StringValue(XEventServiceConstants.XZeroMQ)] XZeroMQ, [StringValue(XEventServiceConstants.XRabbitMQ)] XRabbitMQ, } } ### FILE: C:\Users\SaherElm\Documents\Projects\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 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 = "*"; } } ### FILE: C:\Users\SaherElm\Documents\Projects\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\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Constants\XZeroMQPattern.cs using xExceptions.Attributes; namespace xEventService.Constants { public enum XZeroMQPattern { [StringValue(XEventServiceConstants.None)] None, [StringValue(XEventServiceConstants.PubSub)] PubSub, [StringValue(XEventServiceConstants.PushPull)] PushPull } } ### FILE: C:\Users\SaherElm\Documents\Projects\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 xEventService.Interfaces; using xEventService.Providers; using xExceptions.Constants; namespace xEventService.DI { public static partial class XDIHelperExtension { /// /// Extract EventService Configurations ... /// /// /// public static XEventServiceConfiguration GetXEventServiceConfiguration( 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 AddXEventServiceConfiguration( this IServiceCollection services, IConfiguration configuration ) { // var config = configuration.GetXEventServiceConfiguration(); services.AddXEventServiceConfiguration(config); } /// /// Register Event Service Configurations ... /// /// /// public static void AddXEventServiceConfiguration( this IServiceCollection services, XEventServiceConfiguration configuration ) { // // Validate ... var isValid = configuration.IsValid() && (configuration.Broker == XEventBroker.XKafka || (configuration.Broker == XEventBroker.XZeroMQ && configuration.ZeroMQ.IsValid()) || (configuration.Broker == XEventBroker.XRabbitMQ && configuration.RabbitMQ.IsValid())); if (!isValid) { // Log("Registration Failed, due Invalid Configuration ..."); XException.InvalidConfiguration.Throw(); } // // Register Configuration as Singleton ... services.AddSingleton(configuration); } /// /// Register Service ... /// /// /// /// public static void AddXEventService( this IServiceCollection source, IConfiguration configuration, ServiceLifetime lifeTime = ServiceLifetime.Scoped ) { // var config = configuration.GetXEventServiceConfiguration(); source.AddEventService( lifeTime: lifeTime, configuration: config ); } /// /// Register Service ... /// /// /// /// public static void AddXEventService( this IServiceCollection source, XEventServiceConfiguration configuration, ServiceLifetime lifeTime = ServiceLifetime.Scoped ) { // source.AddEventService( lifeTime: lifeTime, configuration: configuration ); } /// /// Event Consumer Service Registration ... /// /// /// public static void AddXEventServiceConsumer( this IServiceCollection source ) where TConsumer : XBaseEventConsumerBackgroundService { // var config = source.GetRegisteredService(); if (config.IsNullOrDefault()) { // Log("Consumer Registration Faild, due Invalid Configuration issue ..."); XException.InvalidConfiguration.Throw(); } // var consumerType = typeof(TConsumer); switch (config.Broker) { // case XEventBroker.XKafka: // if (!typeof(XKafkaConsumerServiceBase).IsAssignableFrom(consumerType)) { // Log($"Consumer Registration Failed. TConsumer must inherit from XKafkaConsumerServiceBase when Broker is {config.Broker.GetStringValue()} ..."); XException.InvalidConfiguration.Throw(); } break; // case XEventBroker.XZeroMQ: if (!typeof(XZeroMQConsumerServiceBase).IsAssignableFrom(consumerType)) { Log($"Consumer Registration Failed. TConsumer must inherit from XZeroMQConsumerServiceBase when Broker is {config.Broker.GetStringValue()} ..."); XException.InvalidConfiguration.Throw(); } break; // case XEventBroker.XRabbitMQ: if (!typeof(XRabbitMQConsumerServiceBase).IsAssignableFrom(consumerType)) { Log($"Consumer Registration Failed. TConsumer must inherit from XRabbitMQConsumerServiceBase when Broker is {config.Broker.GetStringValue()} ..."); XException.InvalidConfiguration.Throw(); } break; } // source.AddXBackgroundService(); // Log($"Consumer Registered Successfully: {consumerType.Name} for Broker: {config.Broker.GetStringValue()} ..."); } // #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}"); } /// /// Register Event Service ... /// /// /// /// private static void AddEventService( this IServiceCollection source, XEventServiceConfiguration configuration, ServiceLifetime lifeTime = ServiceLifetime.Scoped ) { // // Register Configuration ... source.AddXEventServiceConfiguration(configuration); // // Register based on Broker ... switch (configuration.Broker) { // case XEventBroker.XKafka: source.Add( new ServiceDescriptor(typeof(IXKafkaProducerService), typeof(XKafkaProducerService), lifeTime) ); break; // case XEventBroker.XZeroMQ: source.Add( new ServiceDescriptor(typeof(IXZeroMQProducerService), typeof(XZeroMQProducerService), lifeTime) ); break; // case XEventBroker.XRabbitMQ: source.Add( new ServiceDescriptor(typeof(IXRabbitMQProducerService), typeof(XRabbitMQProducerService), lifeTime) ); break; } // source.Add( new ServiceDescriptor(typeof(IXEventServiceProvider), typeof(XEventServiceProvider), lifeTime) ); // Log($"Service Registered Sucessfull, Provider: {configuration.Broker.GetStringValue()}"); } #endregion } } ### FILE: C:\Users\SaherElm\Documents\Projects\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 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 ... /// /// /// 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; } /// /// 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 ... /// /// /// 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\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXEventServiceProvider.cs namespace xEventService.Interfaces { public interface IXEventServiceProvider : IXProducerServiceBase { } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXKafkaProducerService.cs namespace xEventService.Interfaces { /// /// a Service for Produce Kafka Messages ... /// public interface IXKafkaProducerService : IXProducerServiceBase { } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXProducerServiceBase.cs 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 ); } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXRabbitMQProducerService.cs using System.Threading; using System.Threading.Tasks; using xEventService.Models; namespace xEventService.Interfaces { /// /// RabbitMQ Producer Service Implementation ... /// public interface IXRabbitMQProducerService : IXProducerServiceBase { /// /// 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\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXZeroMQProducerService.cs namespace xEventService.Interfaces { public interface IXZeroMQProducerService : IXProducerServiceBase { } } ### FILE: C:\Users\SaherElm\Documents\Projects\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\SaherElm\Documents\Projects\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\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Models\XRabbitMQOptions.cs 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\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Models\XZeroMQOptions.cs 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; } } ### FILE: C:\Users\SaherElm\Documents\Projects\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 ... /// protected readonly string[] allowedSenders; /// /// Allowed Actions ... /// protected 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(serviceProvider, logger) { // // 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\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Providers\XBaseEventProducerService.cs 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(); } } } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Providers\XEventServiceProvider.cs using System.Threading; using System.Threading.Tasks; 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 XEventServiceProvider : IXEventServiceProvider { // private readonly XEventServiceConfiguration configuration; private readonly IXKafkaProducerService kafkaProducerService; private readonly IXZeroMQProducerService zeroMQProducerService; private readonly IXRabbitMQProducerService rabbitMQProducerService; public XEventServiceProvider( XEventServiceConfiguration configuration, IXKafkaProducerService kafkaProducerService = null, IXZeroMQProducerService zeroMQProducerService = null, IXRabbitMQProducerService rabbitMQProducerService = null ) { // this.configuration = configuration; this.kafkaProducerService = kafkaProducerService; this.zeroMQProducerService = zeroMQProducerService; this.rabbitMQProducerService = rabbitMQProducerService; // // Validate ... var isValid = configuration.IsValid() && configuration.Broker == XEventBroker.XKafka ? kafkaProducerService != null : configuration.Broker == XEventBroker.XZeroMQ ? zeroMQProducerService != null : configuration.Broker == XEventBroker.XRabbitMQ ? rabbitMQProducerService != null : false; if (!isValid) { XException.InvalidConfiguration.Throw(); } } /// /// Send Message to Broker Using XEventRequest instance ... /// /// /// /// public async Task SendMessageAsync( XEventRequest request, CancellationToken cancellationToken = default ) { // // Validate ... if (!request.IsValid()) { XException.InvalidArgs.Throw(); } // var result = await SendMessageAsync( topic: request.Topic, message: request.Message, cancellationToken: cancellationToken ); // return result; } /// /// Send Message to Broker Using XMessage instance ... /// /// /// /// /// public async Task SendMessageAsync( XEventMessage message, string topic = XEventServiceConstants.XDefaultTopic, CancellationToken cancellationToken = default ) { // // Validate Args ... if (!message.IsValid()) { XException.InvalidArgs.Throw(); } // var result = false; switch (configuration.Broker) { // case XEventBroker.XKafka: result = await kafkaProducerService.SendMessageAsync( topic: topic, message: message, cancellationToken: cancellationToken ); break; // case XEventBroker.XZeroMQ: result = await zeroMQProducerService.SendMessageAsync( topic: topic, message: message, cancellationToken: cancellationToken ); break; // case XEventBroker.XRabbitMQ: result = await rabbitMQProducerService.SendMessageAsync( topic: topic, message: message, cancellationToken: cancellationToken ); break; } // return result; } /// /// Dispose ... /// public void Dispose() { // if (kafkaProducerService != null) { kafkaProducerService.Dispose(); } // if (zeroMQProducerService != null) { zeroMQProducerService.Dispose(); } // if (rabbitMQProducerService != null) { rabbitMQProducerService.Dispose(); } } } } ### FILE: C:\Users\SaherElm\Documents\Projects\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 xEventService.Extensions; using xEventService.Models; 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() ) { // var message = consumeResult .Message .Value .FromJSON(); if (message.IsValid( allowedSenders: allowedSenders, allowedActions: allowedActions) ) { // await DoWorkAsync( payload: message, cancellationToken: cancellationToken ); } } } catch (ConsumeException) { await Task.Delay(1000, cancellationToken); } catch (Exception) { await Task.Delay(1000, cancellationToken); } } } /// /// Dispose Object ... /// public override void Dispose() { // if (consumer != null) { consumer.Close(); consumer.Dispose(); } // GC.SuppressFinalize(this); base.Dispose(); } } } ### FILE: C:\Users\SaherElm\Documents\Projects\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 : XBaseEventProducerService, IXKafkaProducerService { private readonly IProducer producer; public XKafkaProducerService( ILogger logger, XEventServiceConfiguration configuration ) : base(configuration, logger) { // // 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() { // if (producer != null) { producer.Dispose(); } // GC.SuppressFinalize(this); } } } ### FILE: C:\Users\SaherElm\Documents\Projects\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(topic); 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() && message.IsValid( allowedSenders: allowedSenders, allowedActions: allowedActions) ) { // 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) { // if (channel != null) { // channel.Close(); channel.Dispose(); } // if (connection != null) { // connection.Close(); connection.Dispose(); } // await base.StopAsync(cancellationToken); } /// /// Dispose Implementation ... /// public override void Dispose() { // if (channel != null) { channel.Dispose(); } // if (connection != null) { 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(string topic) { // channel.ExchangeDeclare( arguments: null, exchange: topic ?? 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: topic ?? configuration.RabbitMQ.Exchange, routingKey: configuration.RabbitMQ.RoutingKey ); } #endregion } } ### FILE: C:\Users\SaherElm\Documents\Projects\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 : 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 { 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 } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Providers\XZeroMQConsumerServiceBase.cs 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.ZeroMQ.IsValid() || configuration.Broker != XEventBroker.XZeroMQ) { 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) { // if (socket != null) { 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 } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Providers\XZeroMQProducerService.cs 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 } }