From 079048dff51ac04e75a654279f459ad1cf937982 Mon Sep 17 00:00:00 2001 From: Hadi Khazaee Asl Date: Wed, 19 Aug 2026 22:44:39 +0330 Subject: [PATCH] Complete xEventService ... --- DI/XDIHelperExtension.cs | 169 +++++++++++++++++++++- Interfaces/IXEventServiceProvider.cs | 5 + Providers/XEventServiceProvider.cs | 159 ++++++++++++++++++++ Providers/XKafkaConsumerServiceBase.cs | 11 +- Providers/XKafkaProducerService.cs | 6 +- Providers/XRabbitMQConsumerServiceBase.cs | 33 +++-- Providers/XRabbitMQProducerService.cs | 8 +- Providers/XZeroMQConsumerServiceBase.cs | 25 +++- 8 files changed, 389 insertions(+), 27 deletions(-) create mode 100644 Interfaces/IXEventServiceProvider.cs create mode 100644 Providers/XEventServiceProvider.cs diff --git a/DI/XDIHelperExtension.cs b/DI/XDIHelperExtension.cs index eba5f25..de4607a 100644 --- a/DI/XDIHelperExtension.cs +++ b/DI/XDIHelperExtension.cs @@ -5,6 +5,8 @@ using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xEventService.Extensions; +using xEventService.Interfaces; +using xEventService.Providers; using xExceptions.Constants; namespace xEventService.DI @@ -16,7 +18,7 @@ namespace xEventService.DI /// /// /// - public static XEventServiceConfiguration GetXEventServiceConfigurations( + public static XEventServiceConfiguration GetXEventServiceConfiguration( this IConfiguration source ) { @@ -34,14 +36,14 @@ namespace xEventService.DI /// /// /// - public static void AddXEventServiceConfigurations( + public static void AddXEventServiceConfiguration( this IServiceCollection services, IConfiguration configuration ) { // - var config = configuration.GetXEventServiceConfigurations(); - services.AddXEventServiceConfigurations(config); + var config = configuration.GetXEventServiceConfiguration(); + services.AddXEventServiceConfiguration(config); } /// @@ -49,14 +51,20 @@ namespace xEventService.DI /// /// /// - public static void AddXEventServiceConfigurations( + public static void AddXEventServiceConfiguration( this IServiceCollection services, XEventServiceConfiguration configuration ) { // // Validate ... - if (!configuration.IsValid()) + 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 ..."); @@ -68,6 +76,104 @@ namespace xEventService.DI 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 ... /// @@ -83,6 +189,57 @@ namespace xEventService.DI { 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 } } \ No newline at end of file diff --git a/Interfaces/IXEventServiceProvider.cs b/Interfaces/IXEventServiceProvider.cs new file mode 100644 index 0000000..7afb179 --- /dev/null +++ b/Interfaces/IXEventServiceProvider.cs @@ -0,0 +1,5 @@ +namespace xEventService.Interfaces +{ + public interface IXEventServiceProvider : IXProducerServiceBase + { } +} \ No newline at end of file diff --git a/Providers/XEventServiceProvider.cs b/Providers/XEventServiceProvider.cs new file mode 100644 index 0000000..80f7fa0 --- /dev/null +++ b/Providers/XEventServiceProvider.cs @@ -0,0 +1,159 @@ +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(); + } + } + } +} \ No newline at end of file diff --git a/Providers/XKafkaConsumerServiceBase.cs b/Providers/XKafkaConsumerServiceBase.cs index 62fd8c6..875e22e 100644 --- a/Providers/XKafkaConsumerServiceBase.cs +++ b/Providers/XKafkaConsumerServiceBase.cs @@ -74,6 +74,7 @@ namespace xEventService.Providers // // Subscribe to Topic ... consumer.Subscribe(topic); + logger.LogInformation($"Kafka Cosume Topic: {topic} ..."); // // Consuming ... @@ -100,18 +101,22 @@ namespace xEventService.Providers { // await DoWorkAsync( - payload: consumeResult.Message, + payload: message, cancellationToken: cancellationToken ); } } } - catch (ConsumeException) + catch (ConsumeException ex) { + // + logger.LogError($"Kafka Consumer Result Parsing Error: {ex.Message} ..."); await Task.Delay(1000, cancellationToken); } - catch (Exception) + catch (Exception ex) { + // + logger.LogError($"Kafka Consumer Result Parsing Error: {ex.Message} ..."); await Task.Delay(1000, cancellationToken); } } diff --git a/Providers/XKafkaProducerService.cs b/Providers/XKafkaProducerService.cs index caccc35..f06963e 100644 --- a/Providers/XKafkaProducerService.cs +++ b/Providers/XKafkaProducerService.cs @@ -63,8 +63,8 @@ namespace xEventService.Providers // result = await SendMessageAsync( - message: request.Message, topic: request.Topic, + message: request.Message, cancellationToken: cancellationToken ); @@ -111,8 +111,10 @@ namespace xEventService.Providers // result = !response.IsNullOrDefault(); } - catch + catch (Exception ex) { + // + logger.LogError($"Kafka Produce Error: {ex.Message} ..."); result = false; } diff --git a/Providers/XRabbitMQConsumerServiceBase.cs b/Providers/XRabbitMQConsumerServiceBase.cs index 5d27e10..b093348 100644 --- a/Providers/XRabbitMQConsumerServiceBase.cs +++ b/Providers/XRabbitMQConsumerServiceBase.cs @@ -51,7 +51,7 @@ namespace xEventService.Providers // SetupConnection(); - DeclareTopology(); + DeclareTopology(topic); consumer = new AsyncEventingBasicConsumer(channel); } @@ -68,6 +68,7 @@ namespace xEventService.Providers // For Fix Blocking Synchronous ... await Task.Yield(); using var cts = new CancellationTokenSource(); + logger.LogInformation($"RabbitMQ Cosume Topic: {topic} ..."); // // Subscribe to Topic ... @@ -103,9 +104,10 @@ namespace xEventService.Providers ); } } - catch + catch (Exception ex) { // + logger.LogError($"RabbitMQ Cosume Error: {ex.Message} ..."); if (configuration.RabbitMQ.ManualAck) { // @@ -153,10 +155,20 @@ namespace xEventService.Providers public override async Task StopAsync(CancellationToken cancellationToken) { // - channel.Close(); - channel.Dispose(); - connection.Close(); - connection.Dispose(); + if (channel != null) + { + // + channel.Close(); + channel.Dispose(); + } + + // + if (connection != null) + { + // + connection.Close(); + connection.Dispose(); + } // await base.StopAsync(cancellationToken); @@ -199,8 +211,9 @@ namespace xEventService.Providers if (result.IsNullOrEmpty()) { // + var exchangeName = topic ?? configuration.RabbitMQ.Exchange; var digits = Guid.NewGuid().ToString().GetDigits(); - result = $"{configuration.RabbitMQ.Exchange}[{digits}]"; + result = $"{exchangeName}[{digits}]"; } // @@ -228,12 +241,12 @@ namespace xEventService.Providers channel = connection.CreateModel(); } - private void DeclareTopology() + private void DeclareTopology(string topic) { // channel.ExchangeDeclare( arguments: null, - exchange: configuration.RabbitMQ.Exchange, + exchange: topic ?? configuration.RabbitMQ.Exchange, durable: configuration.RabbitMQ.ExchangeDurable, autoDelete: configuration.RabbitMQ.ExchangeAutoDelete, type: configuration.RabbitMQ.ExchangeType.GetStringValue() @@ -255,7 +268,7 @@ namespace xEventService.Providers channel.QueueBind( arguments: null, queue: queueName, - exchange: configuration.RabbitMQ.Exchange, + exchange: topic ?? configuration.RabbitMQ.Exchange, routingKey: configuration.RabbitMQ.RoutingKey ); } diff --git a/Providers/XRabbitMQProducerService.cs b/Providers/XRabbitMQProducerService.cs index b462a28..eb83820 100644 --- a/Providers/XRabbitMQProducerService.cs +++ b/Providers/XRabbitMQProducerService.cs @@ -143,7 +143,7 @@ namespace xEventService.Providers { // // Ensure Connection is Alive ... - EnsureChannelAsync(); + EnsureChannel(); // // Serialize Message ... @@ -181,8 +181,10 @@ namespace xEventService.Providers // result = true; } - catch + catch (Exception ex) { + // + logger.LogError($"RabbitMQ Produce Error: {ex.Message} ..."); result = false; } @@ -268,7 +270,7 @@ namespace xEventService.Providers /// /// Ensure Channel is Open and Exchange/Queue Declared ... /// - private void EnsureChannelAsync() + private void EnsureChannel() { // if (IsConnected()) diff --git a/Providers/XZeroMQConsumerServiceBase.cs b/Providers/XZeroMQConsumerServiceBase.cs index 144563a..763f96f 100644 --- a/Providers/XZeroMQConsumerServiceBase.cs +++ b/Providers/XZeroMQConsumerServiceBase.cs @@ -40,7 +40,8 @@ namespace xEventService.Providers { // // Validating Configuration ... - if (configuration.Broker != XEventBroker.XKafka) + if (!configuration.ZeroMQ.IsValid() || + configuration.Broker != XEventBroker.XZeroMQ) { XException.InvalidConfiguration.Throw(); } @@ -61,6 +62,7 @@ namespace xEventService.Providers { // await Task.Yield(); + logger.LogInformation($"ZeroMQ Cosume Topic: {topic} ..."); while (!cancellationToken.IsCancellationRequested) { // @@ -70,6 +72,7 @@ namespace xEventService.Providers string receivedTopic = string.Empty; byte[] body = null; + // if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub) { // @@ -80,6 +83,10 @@ namespace xEventService.Providers { // Process } + else + { + logger.LogError($"ZeroMQ Cosume Error ..."); + } } } else @@ -88,8 +95,13 @@ namespace xEventService.Providers { // Process } + else + { + logger.LogError($"ZeroMQ Cosume Error ..."); + } } + // if (body != null && body.Length > 0) { // @@ -113,8 +125,10 @@ namespace xEventService.Providers await Task.Delay(10, cancellationToken); } } - catch (Exception) + catch (Exception ex) { + // + logger.LogError($"ZeroMQ Cosume Error: {ex.Message} ..."); await Task.Delay(1000, cancellationToken); } } @@ -123,7 +137,12 @@ namespace xEventService.Providers public override async Task StopAsync(CancellationToken cancellationToken) { // - socket.Close(); + if (socket != null) + { + socket.Close(); + } + + // await base.StopAsync(cancellationToken); }