diff --git a/Base/XBaseEvent.cs b/Base/XBaseEvent.cs
deleted file mode 100644
index 2d6c416..0000000
--- a/Base/XBaseEvent.cs
+++ /dev/null
@@ -1,14 +0,0 @@
-using System;
-
-namespace xEventService.Base {
- public abstract class XBaseEvent {
- public string Type { get; protected set; }
- public DateTime Timestamp { get; protected set; }
-
- protected XBaseEvent () {
- //
- Type = GetType ().Name;
- Timestamp = DateTime.UtcNow;
- }
- }
-}
\ No newline at end of file
diff --git a/Configuration/XEventServiceConfiguration.cs b/Configuration/XEventServiceConfiguration.cs
index a1e822e..157dbac 100644
--- a/Configuration/XEventServiceConfiguration.cs
+++ b/Configuration/XEventServiceConfiguration.cs
@@ -1,51 +1,28 @@
using System.Collections.Generic;
using RabbitMQ.Client;
+using xEventService.Constants;
-namespace xEventService.Configuration {
+namespace xEventService.Configuration
+{
///
/// a configuration class for describing how to connect rabbit mq server
/// and how channels must be declared ...
///
- public class XEventServiceConfiguration {
- public string Server { get; set; }
- public int Port { get; set; }
- public string Username { get; set; }
- public string Password { get; set; }
+ public class XEventServiceConfiguration
+ {
+ ///
+ /// Message Broker Address ...
+ ///
+ public string Url { get; set; }
- //
- public bool AllowDynamicQueues { get; set; } = true;
- public XQueueArgs DefaultQueueArgs { get; set; } = new XQueueArgs ();
- public IDictionary Queues { get; set; }
+ ///
+ /// Max Retries to Sending Message ...
+ ///
+ public int MaxRetries { get; set; } = 3;
- //
- public bool AllowDynamicExchangess { get; set; } = true;
- public XExchangeArgs DefaultExchangeArgs { get; set; } = new XExchangeArgs ();
- public IDictionary Exchanges { get; set; }
-
- //
- public static XEventServiceConfiguration GetDefaults () {
- return new XEventServiceConfiguration {
- Port = 5672,
- Username = "guest",
- Password = "guest",
- Server = "localhost",
- Queues = new Dictionary { { "XMainQueue", new XQueueArgs { Name = "XMainQueue" } }
- }
- };
- }
- }
-
- public class XQueueArgs {
- public string Name { get; set; }
- public bool Durable { get; set; } = false;
- public bool Exclusive { get; set; } = false;
- public bool AutoDelete { get; set; } = true;
- }
-
- public class XExchangeArgs {
- public string Name { get; set; }
- public string Type { get; set; } = ExchangeType.Fanout;
- public bool Durable { get; set; } = false;
- public bool AutoDelete { get; set; } = true;
+ ///
+ /// Allowed Message Broker ...
+ ///
+ public XEventBroker Broker { get; set; } = XEventBroker.None;
}
}
\ No newline at end of file
diff --git a/Constants/XEventBroker.cs b/Constants/XEventBroker.cs
new file mode 100644
index 0000000..0d39d85
--- /dev/null
+++ b/Constants/XEventBroker.cs
@@ -0,0 +1,19 @@
+using xExceptions.Attributes;
+
+namespace xEventService.Constants
+{
+ ///
+ /// Available Brokers ...
+ ///
+ public enum XEventBroker
+ {
+ [StringValue("None")]
+ None,
+
+ [StringValue("XKafka")]
+ XKafka,
+
+ [StringValue("XRabbitMQ")]
+ XRabbitMQ
+ }
+}
\ No newline at end of file
diff --git a/Constants/XEventServiceConstants.cs b/Constants/XEventServiceConstants.cs
new file mode 100644
index 0000000..497ed1e
--- /dev/null
+++ b/Constants/XEventServiceConstants.cs
@@ -0,0 +1,13 @@
+namespace xEventService.Constants
+{
+ public struct XEventServiceConstants
+ {
+ //
+ // Service Extensions Log Tag ...
+ public const string XEventServiceDILogTag = "XEventService";
+
+ //
+ // Defaults ...
+ public const string XDefaultTopic = "XSaherelm";
+ }
+}
\ No newline at end of file
diff --git a/DI/XDIHelperExtension.cs b/DI/XDIHelperExtension.cs
index 3db10b9..eba5f25 100644
--- a/DI/XDIHelperExtension.cs
+++ b/DI/XDIHelperExtension.cs
@@ -1,26 +1,29 @@
using System;
-using Microsoft.AspNetCore.Builder;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
-using xEventService.Interfaces;
-using xEventService.Models;
-using xEventService.Providers;
+using xEventService.Extensions;
+using xExceptions.Constants;
-namespace xEventService.DI {
- public static partial class XDIHelperExtension {
+namespace xEventService.DI
+{
+ public static partial class XDIHelperExtension
+ {
///
/// Extract EventService Configurations ...
///
///
///
- public static XEventServiceConfiguration GetXEventServiceConfigurations (this IConfiguration source) {
+ public static XEventServiceConfiguration GetXEventServiceConfigurations(
+ this IConfiguration source
+ )
+ {
//
var configSection = source
- .GetSection (ConfigurationNodeNames.EVENT_SERVICE_NODE_NAME);
- var result = configSection.Get ();
+ .GetSection(ConfigurationNodeNames.EVENT_SERVICE_NODE_NAME);
+ var result = configSection.Get();
//
return result;
@@ -31,13 +34,14 @@ namespace xEventService.DI {
///
///
///
- public static void AddXEventServiceConfigurations (
+ public static void AddXEventServiceConfigurations(
this IServiceCollection services,
IConfiguration configuration
- ) {
+ )
+ {
//
- var config = configuration.GetXEventServiceConfigurations ();
- services.AddXEventServiceConfigurations (config);
+ var config = configuration.GetXEventServiceConfigurations();
+ services.AddXEventServiceConfigurations(config);
}
///
@@ -45,96 +49,40 @@ namespace xEventService.DI {
///
///
///
- public static void AddXEventServiceConfigurations (
+ public static void AddXEventServiceConfigurations(
this IServiceCollection services,
XEventServiceConfiguration configuration
- ) {
+ )
+ {
//
- services.AddSingleton (configuration);
- }
-
- ///
- /// Register Event Service ...
- ///
- ///
- ///
- public static void AddXEventService (
- this IServiceCollection services,
- IConfiguration configuration
- ) {
- //
- // Check Configuration Registered ...
- var config = services.GetRegisteredService ();
- if (config.IsNull ()) {
+ // Validate ...
+ if (!configuration.IsValid())
+ {
//
- Console.WriteLine ($"XEventService: there is no provided XEventServiceConfiguration, using default ...");
-
- //
- // Register Default Configuration ...
- services.AddXEventServiceConfigurations (configuration);
+ Log("Registration Failed, due Invalid Configuration ...");
+ XException.InvalidConfiguration.Throw();
}
//
- // Register IXEventBus ...
- var serviceBas = services.GetRegisteredService ();
- if (serviceBas.IsNull ()) {
- services.AddSingleton ();
- }
+ // Register Configuration as Singleton ...
+ services.AddSingleton(configuration);
}
+ //
+ #region Private ...
+ ///
+ /// a LogTag for Service ...
+ ///
+ private static string XLogTag = XEventServiceConstants.XEventServiceDILogTag;
+
///
- /// Register an Event Handler ...
+ /// print a log in Console ...
///
- ///
- ///
- ///
- ///
- ///
- public static void AddXEventHandler (
- this IServiceCollection services,
- ServiceLifetime lifetime = ServiceLifetime.Transient
- )
- where TEvent : XEvent
- where THandler : IXEventHandler {
- //
- services.Add (new ServiceDescriptor (
- serviceType: typeof (THandler),
- implementationType: typeof (THandler),
- lifetime: lifetime
- ));
-
- //
- services.Add (new ServiceDescriptor (
- serviceType: typeof (IXEventHandler),
- implementationType: typeof (THandler),
- lifetime: lifetime
- ));
- }
-
- ///
- /// Subscribe to specific Event ...
- ///
- ///
- ///
- ///
- ///
- public static void XEventSubscribe (
- this IApplicationBuilder app,
- bool toExchange = true
- )
- where TEvent : XEvent
- where THandler : IXEventHandler {
- //
- var eventName = typeof (TEvent).Name;
-
- //
- // Retrieve EventBus from IOC ...
- var eventBus = app.ApplicationServices.GetRequiredService ();
-
- //
- // Subscribe to Event Handler ...
- var isSubscribed = eventBus.Subscribe ();
- Console.WriteLine ($"XEventService: Event Handler subscription result: {eventName} => {isSubscribed} ...");
+ ///
+ private static void Log(string message)
+ {
+ Console.WriteLine($"{XLogTag} => {message}");
}
+ #endregion
}
}
\ No newline at end of file
diff --git a/Extensions/ConfigurationExtensions.cs b/Extensions/ConfigurationExtensions.cs
deleted file mode 100644
index f557422..0000000
--- a/Extensions/ConfigurationExtensions.cs
+++ /dev/null
@@ -1,288 +0,0 @@
-using System;
-using RabbitMQ.Client;
-using xCommons.Extensions;
-using xEventService.Configuration;
-using xEventService.Models;
-using xExceptions.Constants;
-
-namespace xEventService.Extensions {
- public static class ConfigurationExtensions {
- ///
- /// Create Connection Factory Based on
- /// provided Configuration ...
- ///
- ///
- ///
- public static ConnectionFactory GetConnectionFactory (this XEventServiceConfiguration source) {
- //
- // Validate Args ...
- if (source.IsNull ()) {
- XException.InvalidConfiguration.Throw ();
- }
-
- //
- // Create Connection Factory ...
- var result = new ConnectionFactory {
- Port = source.Port,
- HostName = source.Server,
- UserName = source.Username,
- Password = source.Password
- };
-
- //
- // Return Result ...
- return result;
- }
-
- ///
- /// Create a Connection to Server ...
- ///
- ///
- ///
- public static IConnection Connect (this XEventServiceConfiguration source) {
- //
- // Open Connection ...
- IConnection result = null;
- try {
- //
- // Retrieve Connection Factory ...
- var connectionFactory = source.GetConnectionFactory ();
-
- //
- // Create Connection ...
- result = connectionFactory.CreateConnection ();
- } catch (Exception ex) {
- //
- Console.WriteLine ($"XEventService Exception: Connection Failed, {ex.Message} ...");
-
- //
- XException.ActionFailed.Throw ();
- }
-
- //
- return result;
- }
-
- ///
- /// Validate a Queue Name based on Configurations ...
- ///
- ///
- ///
- public static void ValidateQueue (
- this XEventServiceConfiguration source,
- string queueName
- ) {
- //
- // Validate Args ...
- if (source.IsNull ()) {
- XException.InvalidArgs.Throw ();
- }
-
- //
- // Validate Queues ...
- if ((source.Queues.IsNull () &&
- !source.AllowDynamicQueues) ||
- (!source.Queues.IsNull () &&
- !source.Queues.Keys.HasChild () &&
- !source.AllowDynamicQueues)) {
- //
- Console.WriteLine ($"XEventService Exception: there isn't any configured Queue ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
-
- //
- // Check queueName not empty ...
- if (queueName.IsNullOrEmpty ()) {
- //
- Console.WriteLine ($"XEventService Exception: Invalid QueueName ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
-
- //
- // Check queueName Exists ...
- if (!source.Queues.IsNull () &&
- !source.Queues.ContainsKey (queueName) &&
- !source.AllowDynamicQueues) {
- //
- Console.WriteLine ($"XEventService Exception: QueueName {queueName} not found ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
- }
-
- ///
- /// Validate a Exchange Name based on Configurations ...
- ///
- ///
- ///
- public static void ValidateExchange (
- this XEventServiceConfiguration source,
- string exchangeName
- ) {
- //
- // Validate Args ...
- if (source.IsNull ()) {
- XException.InvalidArgs.Throw ();
- }
-
- //
- // Validate Exchanges ...
- if ((source.Exchanges.IsNull () &&
- !source.AllowDynamicExchangess) ||
- (!source.Exchanges.IsNull () &&
- !source.Exchanges.Keys.HasChild () &&
- !source.AllowDynamicExchangess)) {
- //
- Console.WriteLine ($"XEventService Exception: there isn't any configured Exchanges ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
-
- //
- // Check exchangeName not empty ...
- if (exchangeName.IsNullOrEmpty ()) {
- //
- Console.WriteLine ($"XEventService Exception: Invalid ExchangeName ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
-
- //
- // Check exchangeName Exists ...
- if (!source.Exchanges.IsNull () &&
- !source.Exchanges.ContainsKey (exchangeName) &&
- !source.AllowDynamicExchangess) {
- //
- Console.WriteLine ($"XEventService Exception: ExchangeName {exchangeName} not found ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
- }
-
- ///
- /// Retrieve specific QueueArgs ...
- ///
- ///
- ///
- ///
- public static XQueueArgs GetQueueArgs (
- this XEventServiceConfiguration source,
- string queueName
- ) {
- //
- // Validate Queue Name ...
- source.ValidateQueue (queueName);
-
- //
- var result = source.Queues
- .IsNull () ? null : source
- .Queues[queueName];
- if (result.IsNull ()) {
- //
- // Prepare Default Queue ...
- result = source.DefaultQueueArgs;
- if (result.IsNull ()) {
- result = new XQueueArgs () {
- Name = queueName
- };
- } else {
- result.Name = queueName;
- }
- }
-
- //
- return result;
- }
-
- ///
- /// Retrieve specific ExchangeArgs ...
- ///
- ///
- ///
- ///
- public static XExchangeArgs GetExchangeArgs (
- this XEventServiceConfiguration source,
- string exchangeName
- ) {
- //
- // Validate Exchange Name ...
- source.ValidateExchange (exchangeName);
-
- //
- var result = source.Exchanges
- .IsNull () ? null : source
- .Exchanges[exchangeName];
- if (result.IsNull ()) {
- //
- // Prepare Default Queue ...
- result = source.DefaultExchangeArgs;
- if (result.IsNull ()) {
- result = new XExchangeArgs () {
- Name = exchangeName,
- Type = ExchangeType.Fanout
- };
- } else {
- result.Name = exchangeName;
- }
- }
-
- //
- return result;
- }
-
- ///
- /// Create a Channel ...
- ///
- ///
- ///
- ///
- public static XChannel CreateChannel (
- this XEventServiceConfiguration source,
- string exchangeName = null
- ) {
- //
- // Validate ExchangeName and Retrieve ExchangeArgs if provided ...
- var exchangeArgs = exchangeName
- .IsNullOrEmpty () ?
- null :
- source.GetExchangeArgs (exchangeName);
-
- //
- // Create Connection ...
- var connection = source.Connect ();
- var channel = connection.CreateModel ();
-
- //
- // Create XChannel Model ...
- var result = new XChannel (
- exchangeName: exchangeName,
- connection: connection
- );
-
- //
- // Declaring Queue if Provided ...
- if (!exchangeArgs.IsNull () &&
- !exchangeName.IsNullOrEmpty ()) {
- //
- // Declaring Queue ...
- result.Channel.ExchangeDeclare (
- exchange: exchangeName,
- type: exchangeArgs.Type,
- durable: exchangeArgs.Durable,
- autoDelete: exchangeArgs.AutoDelete
- );
- }
-
- //
- return result;
- }
- }
-}
\ No newline at end of file
diff --git a/Extensions/XEventModelExtensions.cs b/Extensions/XEventModelExtensions.cs
deleted file mode 100644
index 1bf903f..0000000
--- a/Extensions/XEventModelExtensions.cs
+++ /dev/null
@@ -1,44 +0,0 @@
-using System;
-using System.Text;
-using xCommons.Extensions;
-
-namespace xEventService.Extensions {
- public static class XEventModelExtensions {
- ///
- /// Convert an Object to bytes array for publishing
- ///
- ///
- ///
- ///
- public static byte[] ToBody (this T source) {
- //
- var json = source
- .ToJSON ();
-
- //
- var result = json
- .ToBytes ();
-
- //
- return result;
- }
-
- ///
- /// Convert Recieved Bytes to Specific Type ...
- ///
- ///
- ///
- ///
- public static T FromBody (this ReadOnlyMemory source) {
- //
- var jsonString = Encoding.UTF8
- .GetString (
- source.ToArray ()
- );
- var result = jsonString.FromJSON ();
-
- //
- return result;
- }
- }
-}
\ No newline at end of file
diff --git a/Extensions/XEventServiceExtensions.cs b/Extensions/XEventServiceExtensions.cs
new file mode 100644
index 0000000..421a7c8
--- /dev/null
+++ b/Extensions/XEventServiceExtensions.cs
@@ -0,0 +1,81 @@
+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;
+ }
+ }
+}
\ No newline at end of file
diff --git a/Interfaces/IXEventBus.cs b/Interfaces/IXEventBus.cs
deleted file mode 100644
index 547b946..0000000
--- a/Interfaces/IXEventBus.cs
+++ /dev/null
@@ -1,11 +0,0 @@
-using xEventService.Models;
-
-namespace xEventService.Interfaces {
- public interface IXEventBus {
- bool Publish (T @event) where T : XEvent;
-
- bool Subscribe ()
- where TEvent : XEvent
- where THandler : IXEventHandler;
- }
-}
\ No newline at end of file
diff --git a/Interfaces/IXEventHandler.cs b/Interfaces/IXEventHandler.cs
deleted file mode 100644
index 63ec78c..0000000
--- a/Interfaces/IXEventHandler.cs
+++ /dev/null
@@ -1,11 +0,0 @@
-using System.Threading.Tasks;
-using xEventService.Models;
-
-namespace xEventService.Interfaces {
- public interface IXEventHandler : IXEventHandler
- where TEvent : XEvent {
- Task HandleAsync (TEvent @event);
- }
-
- public interface IXEventHandler { }
-}
\ No newline at end of file
diff --git a/Interfaces/IXKafkaProducerService.cs b/Interfaces/IXKafkaProducerService.cs
new file mode 100644
index 0000000..da8e0ad
--- /dev/null
+++ b/Interfaces/IXKafkaProducerService.cs
@@ -0,0 +1,35 @@
+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,
+ string topic = XEventServiceConstants.XDefaultTopic,
+ CancellationToken cancellationToken = default
+ );
+
+ ///
+ /// Sending a Message ...
+ ///
+ ///
+ ///
+ ///
+ ///
+ Task SendMessageAsync(
+ XEventMessage message,
+ string topic = XEventServiceConstants.XDefaultTopic,
+ CancellationToken cancellationToken = default
+ );
+ }
+}
\ No newline at end of file
diff --git a/Models/XChannel.cs b/Models/XChannel.cs
deleted file mode 100644
index 610f0f8..0000000
--- a/Models/XChannel.cs
+++ /dev/null
@@ -1,224 +0,0 @@
-using System;
-using RabbitMQ.Client;
-using RabbitMQ.Client.Events;
-using xCommons.Extensions;
-using xEventService.Configuration;
-using xEventService.Extensions;
-using xExceptions.Constants;
-
-namespace xEventService.Models {
- public class XChannel {
- //
- public IModel Channel { get; protected set; }
- public string QueueName { get; protected set; }
- public string ExchangeName { get; protected set; }
- public IConnection Connection { get; protected set; }
-
- //
- #region Constructors ...
- public XChannel (
- string exchangeName = null,
- IConnection connection = null,
- IModel channel = null,
- EventingBasicConsumer consumer = null
- ) {
- //
- // Validate Connection ...
- if (connection.IsNull ()) {
- XException.InvalidArgs.Throw ();
- }
- Connection = connection;
-
- //
- // Validate Channel ...
- if (channel.IsNull ()) {
- channel = connection.CreateModel ();
- }
- Channel = channel;
-
- //
- // Check Queue Name and it's Declaration ...
- if (!exchangeName.IsNullOrEmpty ()) {
- //
- ExchangeName = exchangeName;
- }
- }
-
- public XChannel (
- XEventServiceConfiguration configuration,
- bool dispatchConsumersAsync = true,
- string exchangeName = null
- ) {
- //
- // Generate Factory ...
- var factory = configuration.GetConnectionFactory ();
- if (dispatchConsumersAsync) {
- factory.DispatchConsumersAsync = true;
- }
-
- //
- // Create Cahnnel ...
- Connection = factory.CreateConnection ();
- Channel = Connection.CreateModel ();
-
- //
- // Check Exchange Declaration if provided ...
- if (!exchangeName.IsNullOrEmpty ()) {
- //
- ExchangeName = exchangeName;
- var exchangeArgs = configuration.GetExchangeArgs (ExchangeName);
-
- //
- // Declaring Exchange ...
- Channel.ExchangeDeclare (
- exchange: exchangeName,
- type: exchangeArgs.Type,
- durable: exchangeArgs.Durable,
- autoDelete: exchangeArgs.AutoDelete
- );
-
- //
- // Create a Queue ...
- var queueName = $"{exchangeName}[{Guid.NewGuid().ToString().GetDigits()}]";
- var queueArgs = configuration.GetQueueArgs (queueName);
- QueueName = queueName;
- Channel.QueueDeclare (
- queue: queueName,
- durable: queueArgs.Durable,
- exclusive: queueArgs.Exclusive,
- autoDelete: queueArgs.AutoDelete
- );
-
- //
- // Bind Queue to Exchange ...
- Channel.QueueBind (
- queue: queueName,
- exchange: exchangeName,
- routingKey: ""
- );
- }
- }
- #endregion
-
- //
- #region Actions ...
- ///
- /// Close Channel ...
- ///
- public void CloseChannel () {
- //
- // Check Channel Exists ...
- if (!this.Channel.IsNull ()) {
- //
- // Close Channel if it's Open ...
- if (this.Channel.IsOpen) {
- this.Channel.Close ();
- }
-
- //
- this.Channel.Dispose ();
- this.Channel = null;
- }
- }
-
- ///
- /// Close Connection ...
- ///
- public void CloseConnection () {
- //
- // Check Connection Exists ...
- if (!this.Connection.IsNull ()) {
- //
- // Close Connection if it's Open ...
- if (this.Connection.IsOpen) {
- this.Connection.Close ();
- }
-
- //
- this.Connection.Dispose ();
- this.Connection = null;
- }
- }
-
- ///
- /// Publish a Message through Channel on Declare Queue ...
- ///
- ///
- ///
- ///
- public bool Publish (T message) {
- //
- // Validate Args ...
- if (this.Connection.IsNull () ||
- !this.Connection.IsOpen ||
- this.Channel.IsNull () ||
- !this.Channel.IsOpen) {
- //
- Console.WriteLine ($"XEventService Exception: Publish Failed ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
-
- //
- // Check Publish Method ...
- //
- // Check Exchange Name ...
- if (this.ExchangeName.IsNullOrEmpty ()) {
- //
- Console.WriteLine ($"XEventService Exception: Publish Failed, Invalid Exchange Name {ExchangeName} ...");
-
- //
- XException.InvalidArgs.Throw ();
- }
-
- //
- try {
- //
- // Publish Message through Channel to Declared Queue ...
- var body = message.ToBody ();
- this.Channel.BasicPublish (
- exchange: ExchangeName,
- routingKey: "",
- basicProperties : null,
- body : body
- );
-
- //
- return true;
- } catch (Exception ex) {
- //
- Console.WriteLine ($"XEventService Exception: Publish Failed, {ex.Message} ...");
-
- //
- return false;
- }
- }
-
- public void Consume (EventingBasicConsumer consumer) {
- Channel.BasicConsume (
- queue: QueueName,
- autoAck: true,
- consumer: consumer
- );
- }
-
- public void ConsumeAsync (AsyncEventingBasicConsumer consumer) {
- Channel.BasicConsume (
- queue: QueueName,
- autoAck: true,
- consumer: consumer
- );
- }
-
- ///
- /// Dispose XChannel ...
- ///
- public void Dispose () {
- //
- this.CloseChannel ();
- this.CloseConnection ();
- }
- #endregion
- }
-}
\ No newline at end of file
diff --git a/Models/XEvent.cs b/Models/XEvent.cs
deleted file mode 100644
index 2fdde49..0000000
--- a/Models/XEvent.cs
+++ /dev/null
@@ -1,5 +0,0 @@
-using xEventService.Base;
-
-namespace xEventService.Models {
- public class XEvent : XBaseEvent { }
-}
\ No newline at end of file
diff --git a/Models/XEventMessage.cs b/Models/XEventMessage.cs
new file mode 100644
index 0000000..51a69f6
--- /dev/null
+++ b/Models/XEventMessage.cs
@@ -0,0 +1,40 @@
+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();
+ }
+}
\ No newline at end of file
diff --git a/Models/XEventRequest.cs b/Models/XEventRequest.cs
new file mode 100644
index 0000000..ba1a5e0
--- /dev/null
+++ b/Models/XEventRequest.cs
@@ -0,0 +1,18 @@
+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; }
+ }
+}
\ No newline at end of file
diff --git a/Providers/XEventBus.cs b/Providers/XEventBus.cs
deleted file mode 100644
index c0a9334..0000000
--- a/Providers/XEventBus.cs
+++ /dev/null
@@ -1,260 +0,0 @@
-using System;
-using System.Collections.Generic;
-using System.Linq;
-using System.Text;
-using System.Threading.Tasks;
-using Microsoft.Extensions.DependencyInjection;
-using Microsoft.Extensions.Logging;
-using RabbitMQ.Client.Events;
-using xCommons.Extensions;
-using xEventService.Configuration;
-using xEventService.Extensions;
-using xEventService.Interfaces;
-using xEventService.Models;
-using xExceptions.Constants;
-
-namespace xEventService.Providers {
- public sealed class XEventBus : IXEventBus {
- //
- #region Properties ...
- private readonly List events;
- private readonly ILogger logger;
- private readonly IServiceScopeFactory serviceScopeFactory;
- private readonly XEventServiceConfiguration configuration;
- private readonly IDictionary> handlers;
- #endregion
-
- //
- #region Constructors ...
- public XEventBus (
- ILoggerFactory loggerFactory,
- IServiceScopeFactory serviceScopeFactory,
- XEventServiceConfiguration configuration
- ) {
- //
- this.configuration = configuration;
- this.serviceScopeFactory = serviceScopeFactory;
-
- //
- events = new List ();
- handlers = new Dictionary> ();
- logger = loggerFactory.CreateLogger ();
- }
- #endregion
-
- //
- #region Abstracts ...
- public bool Publish (T @event) where T : XEvent {
- //
- // Validate Args ...
- if (@event.IsNull ()) {
- return false;
- }
-
- //
- // Prepare Channel ...
- XChannel channel = null;
-
- //
- // Extract event name ...
- var exchangeName = @event.GetType ().Name;
-
- //
- // Do Action ...
- try {
- //
- // Get XChannel Model ...
- channel = configuration.CreateChannel (exchangeName);
-
- //
- // Publish Command ...
- channel.Publish (message: @event);
-
- //
- // Return Result ...
- return true;
- } catch (Exception ex) {
- //
- logger.LogError ($"XEventService Exception: Publishing Failed {exchangeName}/{ex.Message} ...");
-
- //
- return false;
- } finally {
- //
- if (!channel.IsNull ()) {
- channel.Dispose ();
- }
- }
- }
-
- public bool Subscribe ()
- where TEvent : XEvent
- where THandler : IXEventHandler {
- //
- // Extract required data ...
- var eventType = typeof (TEvent);
- var eventName = eventType.Name;
- var handlerType = typeof (THandler);
-
- //
- // Check event exists in list or not ...
- if (!events.Contains (eventType)) {
- //
- // Add it if not exists ...
- events.Add (eventType);
- }
-
- //
- // Check handlers Dictionary contains list or not ...
- if (!handlers.ContainsKey (eventName)) {
- handlers.Add (eventName, new List ());
- }
-
- //
- // Check handlers Subscribed before or not ...
- var isHandlerSubscribed = handlers[eventName]
- .Any (h => h.GetType () == handlerType);
- if (isHandlerSubscribed) {
- //
- logger.LogError ($"XEventService Exception: Duplicate Event Handler Subscription {eventName}/{handlerType.Name} ...");
-
- //
- return false;
- }
-
- //
- // Register Handler ...
- handlers[eventName].Add (handlerType);
-
- //
- // Do Consuming Handler ...
- ConsumeEvent ();
-
- //
- return true;
- }
- #endregion
-
- //
- #region Private ...
- ///
- /// Consume Specific Event ...
- ///
- ///
- private void ConsumeEvent ()
- where TEvent : XEvent {
- //
- // Extract required data ...
- var eventType = typeof (TEvent);
- var eventName = eventType.Name;
-
- //
- // Create ChannelObject ...
- var channel = new XChannel (
- exchangeName: eventName,
- configuration: configuration,
- dispatchConsumersAsync: true
- );
-
- //
- // Create Consumer ...
- var consumer = new AsyncEventingBasicConsumer (channel.Channel);
-
- //
- // Set Consumer Delegate ...
- consumer.Received += XConsumerReceivedDelegate;
-
- //
- // Consume Async Consumer ...
- channel.ConsumeAsync (consumer);
- }
-
- ///
- /// this Delegate method calls when a message recieved ...
- ///
- private async Task XConsumerReceivedDelegate (
- object sender,
- BasicDeliverEventArgs ea
- ) {
- //
- // Generate Required Data ...
- var eventName = ea.Exchange;
- var json = Encoding.UTF8
- .GetString (
- ea.Body
- .ToArray ()
- );
-
- //
- // Do Processing Event ...
- try {
- //
- // Process Event in non Blocking Tasks ...
- await ProcessEvent (eventName, json)
- .ConfigureAwait (false);
- } catch (Exception ex) {
- //
- logger.LogError ($"XEventService Exception: Processing Event failed {eventName}/{ex.Message} ...");
-
- //
- XException.ActionFailed.Throw ();
- } finally {
- //
- var channel = ((EventingBasicConsumer) sender).Model;
- channel.BasicAck (
- deliveryTag: ea.DeliveryTag,
- multiple: false
- );
- }
- }
-
- ///
- /// Here we must call registered Event Handlers to Handle the Event ...
- ///
- private async Task ProcessEvent (
- string eventName,
- string json
- ) {
- //
- // Check Handler Exists for current Event ...
- var isExistsHandler = handlers.ContainsKey (eventName);
- if (!isExistsHandler) {
- return;
- }
-
- //
- // Using Service Scope Factory ...
- using (var scope = serviceScopeFactory.CreateScope ()) {
- //
- // Get Handler Subscriptions ...
- var subscriptions = handlers[eventName];
- foreach (var subscription in subscriptions) {
- //
- // Inject Handler from Dependency Injections ...
- var handler = scope.ServiceProvider.GetService (subscription);
- if (handler.IsNull ()) {
- //
- logger.LogInformation ($"XEventService: there is no handler Registered in DependencyInjection for {eventName}/{subscription.Name}");
-
- //
- continue;
- }
-
- //
- // Retrieve Event Type ...
- var eventType = events.SingleOrDefault (t => t.Name == eventName);
- var @event = json.FromJSON (eventType);
- var conreteType = typeof (IXEventHandler<>).MakeGenericType (eventType);
-
- //
- // Call EventHandler 'HandleAsync' method ...
- await ((Task) conreteType
- .GetMethod (nameof (IXEventHandler.HandleAsync))
- .Invoke (handler, new object[] { @event }))
- .ConfigureAwait (true);
- }
- }
- }
- #endregion
- }
-}
\ No newline at end of file
diff --git a/Providers/XKafkaProducerService.cs b/Providers/XKafkaProducerService.cs
new file mode 100644
index 0000000..e92339b
--- /dev/null
+++ b/Providers/XKafkaProducerService.cs
@@ -0,0 +1,32 @@
+using xCommons.Extensions;
+using xExceptions.Constants;
+using xEventService.Extensions;
+using xEventService.Interfaces;
+using xEventService.Configuration;
+
+namespace xEventService.Providers
+{
+ public class XKafkaProducerService : IXKafkaProducerService
+ {
+ private readonly IProducer producer;
+ private readonly ILogger logger;
+
+ public XKafkaProducerService(
+ ILoggeer logger,
+ XEventServiceConfiguration configuration
+ )
+ {
+ //
+ this.logger = logger;
+
+ //
+ // Validate Configurations ...
+ if (!configuration.IsValid() ||
+ configuration.Broker != Constants.XEventBroker.XKafka
+ )
+ {
+ XException.InvalidConfiguration.Throw();
+ }
+ }
+ }
+}
\ No newline at end of file
diff --git a/xEventService.csproj b/xEventService.csproj
index dd5af93..ebc3ecc 100644
--- a/xEventService.csproj
+++ b/xEventService.csproj
@@ -36,6 +36,7 @@
+