From 28e679235c5585da9d126d0b063be339f11f4ffd Mon Sep 17 00:00:00 2001 From: Hadi Khazaee Asl Date: Wed, 19 Aug 2026 18:05:51 +0330 Subject: [PATCH] add Register DI Producers ... --- DI/XDIHelperExtension.cs | 77 +++++++++++++ Interfaces/IXEventServiceProvider.cs | 5 + Providers/XEventServiceProvider.cs | 161 +++++++++++++++++++++++++++ 3 files changed, 243 insertions(+) create mode 100644 Interfaces/IXEventServiceProvider.cs create mode 100644 Providers/XEventServiceProvider.cs diff --git a/DI/XDIHelperExtension.cs b/DI/XDIHelperExtension.cs index eba5f25..7126b90 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 @@ -68,6 +70,25 @@ 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 ... /// @@ -83,6 +104,62 @@ 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/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..bf2cc20 --- /dev/null +++ b/Providers/XEventServiceProvider.cs @@ -0,0 +1,161 @@ +using System; +using System.Collections.Generic; +using System.Linq; +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