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