diff --git a/DI/XDIHelperExtension.cs b/DI/XDIHelperExtension.cs index 7126b90..ca4a50d 100644 --- a/DI/XDIHelperExtension.cs +++ b/DI/XDIHelperExtension.cs @@ -18,7 +18,7 @@ namespace xEventService.DI /// /// /// - public static XEventServiceConfiguration GetXEventServiceConfigurations( + public static XEventServiceConfiguration GetXEventServiceConfiguration( this IConfiguration source ) { @@ -36,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); } /// @@ -51,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 ..."); @@ -70,25 +76,6 @@ namespace xEventService.DI services.AddSingleton(configuration); } - /// - /// Register Service ... - /// - /// - /// - /// - public static void AddXEventService( - this IServiceCollection source, - XEventServiceConfiguration configuration, - ServiceLifetime lifeTime = ServiceLifetime.Scoped - ) - { - // - source.AddEventService( - lifeTime: lifeTime, - configuration: configuration - ); - } - // #region Private ... /// @@ -104,62 +91,6 @@ namespace xEventService.DI { Console.WriteLine($"{XLogTag} => {message}"); } - - /// - /// Register Event Service ... - /// - /// - /// - /// - private static void AddEventService( - this IServiceCollection source, - XEventServiceConfiguration configuration, - ServiceLifetime lifeTime = ServiceLifetime.Scoped - ) - { - // - // Validate ... - if (!configuration.IsValid()) - { - // - Log("Initialization Failed due Invalid Configuration ..."); - XException.InvalidConfiguration.Throw(); - } - - // - // 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/Providers/XEventServiceProvider.cs b/Providers/XEventServiceProvider.cs index bf2cc20..80f7fa0 100644 --- a/Providers/XEventServiceProvider.cs +++ b/Providers/XEventServiceProvider.cs @@ -1,6 +1,3 @@ -using System; -using System.Collections.Generic; -using System.Linq; using System.Threading; using System.Threading.Tasks; using xCommons.Extensions; @@ -15,6 +12,7 @@ namespace xEventService.Providers { public class XEventServiceProvider : IXEventServiceProvider { + // private readonly XEventServiceConfiguration configuration; private readonly IXKafkaProducerService kafkaProducerService; private readonly IXZeroMQProducerService zeroMQProducerService; 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); }