### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Configuration\XEventServiceConfiguration.cs using System.Collections.Generic; using RabbitMQ.Client; using xEventService.Constants; 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; } } ### 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("None")] None, [StringValue("XKafka")] XKafka, [StringValue("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"; } } ### 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 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\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xEventService\Extensions\XEventServiceExtensions.cs using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xEventService.Models; 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 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; } } } ### FILE: C:\Users\SaherElm\Documents\Projects\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\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\Providers\XKafkaConsumerServiceBase.cs using System; using Confluent.Kafka; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using xEventService.Configuration; using xEventService.Constants; namespace xEventService.Providers { /// /// an Abstract Kafka Consumer Service ... /// public abstract class XKafkaConsumerServiceBase : BackgroundService { /// /// Consumed Topic Messages ... /// private readonly string topic; /// /// Kafka Consumer Class ... /// public readonly IConsumer consumer; private readonly XEventServiceConfiguration configuration; private readonly ILogger logger; protected XKafkaConsumerServiceBase( XEventServiceConfiguration configuration, ILogger logger, string topic = XEventServiceConstants.XDefaultTopic ) { // this.topic = topic; this.logger = logger; this.configuration = configuration; // // Prepare Kafka Consumer Configuration ... var config = new ConsumerConfig { EnableAutoCommit = false, BootstrapServers = configuration.Url, AutoOffsetReset = AutoOffsetReset.Earliest, GroupId = XEventServiceConstants.XDefaultConsumerGroup, }; consumer = new ConsumerBuilder(config).Build(); } /// /// Dispose Implementation ... /// public override void Dispose() { // 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 : 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 = true; } catch { result = false; } // return result; } /// /// Dispose Required Objects ... /// public void Dispose() { // producer.Dispose(); GC.SuppressFinalize(this); } } }