Compare commits

...
11 Commits
35 changed files with 2327 additions and 997 deletions
-14
View File
@@ -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;
}
}
}
+30 -42
View File
@@ -1,51 +1,39 @@
using System.Collections.Generic;
using RabbitMQ.Client;
using xEventService.Models;
using xEventService.Constants;
namespace xEventService.Configuration {
namespace xEventService.Configuration
{
/// <summary>
/// a configuration class for describing how to connect rabbit mq server
/// and how channels must be declared ...
/// </summary>
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
{
/// <summary>
/// Message Broker Address ...
/// </summary>
public string Url { get; set; }
//
public bool AllowDynamicQueues { get; set; } = true;
public XQueueArgs DefaultQueueArgs { get; set; } = new XQueueArgs ();
public IDictionary<string, XQueueArgs> Queues { get; set; }
/// <summary>
/// Max Retries to Sending Message ...
/// </summary>
public int MaxRetries { get; set; } = 3;
//
public bool AllowDynamicExchangess { get; set; } = true;
public XExchangeArgs DefaultExchangeArgs { get; set; } = new XExchangeArgs ();
public IDictionary<string, XExchangeArgs> Exchanges { get; set; }
//
public static XEventServiceConfiguration GetDefaults () {
return new XEventServiceConfiguration {
Port = 5672,
Username = "guest",
Password = "guest",
Server = "localhost",
Queues = new Dictionary<string, XQueueArgs> { { "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;
/// <summary>
/// Allowed Message Broker ...
/// </summary>
public XEventBroker Broker { get; set; } = XEventBroker.None;
/// <summary>
/// ZeroMQ Specific Options ...
/// Only used when Broker is XZeroMQ ...
/// </summary>
public XZeroMQOptions ZeroMQ { get; set; } = new XZeroMQOptions();
/// <summary>
/// RabbitMQ Specific Options ...
/// Only used when Broker is XRabbitMQ ...
/// </summary>
public XRabbitMQOptions RabbitMQ { get; set; } = new XRabbitMQOptions();
}
}
+22
View File
@@ -0,0 +1,22 @@
using xExceptions.Attributes;
namespace xEventService.Constants
{
/// <summary>
/// Available Brokers ...
/// </summary>
public enum XEventBroker
{
[StringValue(XEventServiceConstants.None)]
None,
[StringValue(XEventServiceConstants.XKafka)]
XKafka,
[StringValue(XEventServiceConstants.XZeroMQ)]
XZeroMQ,
[StringValue(XEventServiceConstants.XRabbitMQ)]
XRabbitMQ,
}
}
+25
View File
@@ -0,0 +1,25 @@
namespace xEventService.Constants
{
public struct XEventServiceConstants
{
//
// Service Extensions Log Tag ...
public const string XEventServiceDILogTag = "XEventService";
//
// Defaults ...
public const string XDefaultTopic = "XSaherelm";
public const string XDefaultConsumerGroup = "XSaherElmGroup";
//
public const string None = "None";
public const string XKafka = "XKafka";
public const string PubSub = "PubSub";
public const string XZeroMQ = "XZeroMQ";
public const string PushPull = "PushPull";
public const string XRabbitMQ = "XRabbitMQ";
//
public const string XWildcard = "*";
}
}
+37
View File
@@ -0,0 +1,37 @@
using xExceptions.Attributes;
namespace xEventService.Constants
{
/// <summary>
/// Available RabbitMQ Exchange Types ...
/// </summary>
public struct XRabbitMQExchangeTypes
{
public const string None = "none";
public const string Direct = "direct";
public const string Fanout = "fanout";
public const string Topic = "topic";
public const string Headers = "headers";
}
/// <summary>
/// RabbitMQ Exchange Type Enum ...
/// </summary>
public enum XRabbitMQExchangeType
{
[StringValue(XRabbitMQExchangeTypes.None)]
None,
[StringValue(XRabbitMQExchangeTypes.Direct)]
Direct,
[StringValue(XRabbitMQExchangeTypes.Fanout)]
Fanout,
[StringValue(XRabbitMQExchangeTypes.Topic)]
Topic,
[StringValue(XRabbitMQExchangeTypes.Headers)]
Headers,
}
}
+16
View File
@@ -0,0 +1,16 @@
using xExceptions.Attributes;
namespace xEventService.Constants
{
public enum XZeroMQPattern
{
[StringValue(XEventServiceConstants.None)]
None,
[StringValue(XEventServiceConstants.PubSub)]
PubSub,
[StringValue(XEventServiceConstants.PushPull)]
PushPull
}
}
+49 -95
View File
@@ -1,26 +1,29 @@
using System;
using Microsoft.AspNetCore.Builder;
using xCommons.Extensions;
using xExceptions.Constants;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Configuration;
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;
namespace xEventService.DI {
public static partial class XDIHelperExtension {
namespace xEventService.DI
{
public static partial class XDIHelperExtension
{
/// <summary>
/// Extract EventService Configurations ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static XEventServiceConfiguration GetXEventServiceConfigurations (this IConfiguration source) {
public static XEventServiceConfiguration GetXEventServiceConfiguration(
this IConfiguration source
)
{
//
var configSection = source
.GetSection (ConfigurationNodeNames.EVENT_SERVICE_NODE_NAME);
var result = configSection.Get<XEventServiceConfiguration> ();
.GetSection(ConfigurationNodeNames.EVENT_SERVICE_NODE_NAME);
var result = configSection.Get<XEventServiceConfiguration>();
//
return result;
@@ -31,13 +34,14 @@ namespace xEventService.DI {
/// </summary>
/// <param name="services"></param>
/// <param name="configuration"></param>
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);
}
/// <summary>
@@ -45,96 +49,46 @@ namespace xEventService.DI {
/// </summary>
/// <param name="services"></param>
/// <param name="configuration"></param>
public static void AddXEventServiceConfigurations (
public static void AddXEventServiceConfiguration(
this IServiceCollection services,
XEventServiceConfiguration configuration
) {
)
{
//
services.AddSingleton<XEventServiceConfiguration> (configuration);
}
/// <summary>
/// Register Event Service ...
/// </summary>
/// <param name="services"></param>
/// <param name="configuration"></param>
public static void AddXEventService (
this IServiceCollection services,
IConfiguration configuration
) {
//
// Check Configuration Registered ...
var config = services.GetRegisteredService<XEventServiceConfiguration> ();
if (config.IsNull ()) {
// Validate ...
var isValid = configuration.IsValid() &&
(configuration.Broker == XEventBroker.XKafka ||
(configuration.Broker == XEventBroker.XZeroMQ &&
configuration.ZeroMQ.IsValid()) ||
(configuration.Broker == XEventBroker.XRabbitMQ &&
configuration.RabbitMQ.IsValid()));
if (!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<IXEventBus> ();
if (serviceBas.IsNull ()) {
services.AddSingleton<IXEventBus, XEventBus> ();
}
// Register Configuration as Singleton ...
services.AddSingleton(configuration);
}
//
#region Private ...
/// <summary>
/// a LogTag for Service ...
/// </summary>
private static string XLogTag = XEventServiceConstants.XEventServiceDILogTag;
/// <summary>
/// Register an Event Handler ...
/// print a log in Console ...
/// </summary>
/// <param name="services"></param>
/// <param name="lifetime"></param>
/// <typeparam name="TEvent"></typeparam>
/// <typeparam name="THandler"></typeparam>
/// <returns></returns>
public static void AddXEventHandler<TEvent, THandler> (
this IServiceCollection services,
ServiceLifetime lifetime = ServiceLifetime.Transient
)
where TEvent : XEvent
where THandler : IXEventHandler<TEvent> {
//
services.Add (new ServiceDescriptor (
serviceType: typeof (THandler),
implementationType: typeof (THandler),
lifetime: lifetime
));
//
services.Add (new ServiceDescriptor (
serviceType: typeof (IXEventHandler<TEvent>),
implementationType: typeof (THandler),
lifetime: lifetime
));
}
/// <summary>
/// Subscribe to specific Event ...
/// </summary>
/// <param name="app"></param>
/// <typeparam name="TEvent"></typeparam>
/// <typeparam name="THandler"></typeparam>
/// <returns></returns>
public static void XEventSubscribe<TEvent, THandler> (
this IApplicationBuilder app,
bool toExchange = true
)
where TEvent : XEvent
where THandler : IXEventHandler<TEvent> {
//
var eventName = typeof (TEvent).Name;
//
// Retrieve EventBus from IOC ...
var eventBus = app.ApplicationServices.GetRequiredService<IXEventBus> ();
//
// Subscribe to Event Handler ...
var isSubscribed = eventBus.Subscribe<TEvent, THandler> ();
Console.WriteLine ($"XEventService: Event Handler subscription result: {eventName} => {isSubscribed} ...");
/// <param name="message"></param>
private static void Log(string message)
{
Console.WriteLine($"{XLogTag} => {message}");
}
#endregion
}
}
-288
View File
@@ -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 {
/// <summary>
/// Create Connection Factory Based on
/// provided Configuration ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
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;
}
/// <summary>
/// Create a Connection to Server ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
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;
}
/// <summary>
/// Validate a Queue Name based on Configurations ...
/// </summary>
/// <param name="source"></param>
/// <param name="queueName"></param>
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 ();
}
}
/// <summary>
/// Validate a Exchange Name based on Configurations ...
/// </summary>
/// <param name="source"></param>
/// <param name="exchangeName"></param>
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 ();
}
}
/// <summary>
/// Retrieve specific QueueArgs ...
/// </summary>
/// <param name="source"></param>
/// <param name="queueName"></param>
/// <returns></returns>
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;
}
/// <summary>
/// Retrieve specific ExchangeArgs ...
/// </summary>
/// <param name="source"></param>
/// <param name="exchangeName"></param>
/// <returns></returns>
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;
}
/// <summary>
/// Create a Channel ...
/// </summary>
/// <param name="source"></param>
/// <param name="exchangeName"></param>
/// <returns></returns>
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;
}
}
}
-44
View File
@@ -1,44 +0,0 @@
using System;
using System.Text;
using xCommons.Extensions;
namespace xEventService.Extensions {
public static class XEventModelExtensions {
/// <summary>
/// Convert an Object to bytes array for publishing
/// </summary>
/// <param name="source"></param>
/// <typeparam name="T"></typeparam>
/// <returns></returns>
public static byte[] ToBody<T> (this T source) {
//
var json = source
.ToJSON ();
//
var result = json
.ToBytes ();
//
return result;
}
/// <summary>
/// Convert Recieved Bytes to Specific Type ...
/// </summary>
/// <param name="source"></param>
/// <typeparam name="T"></typeparam>
/// <returns></returns>
public static T FromBody<T> (this ReadOnlyMemory<byte> source) {
//
var jsonString = Encoding.UTF8
.GetString (
source.ToArray ()
);
var result = jsonString.FromJSON<T> ();
//
return result;
}
}
}
+344
View File
@@ -0,0 +1,344 @@
using System.Linq;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using xEventService.Constants;
using xEventService.Configuration;
namespace xEventService.Extensions
{
/// <summary>
/// Extensions for Event Service ...
/// </summary>
public static class XEventServiceExtensions
{
/// <summary>
/// Validate Event Broker ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XEventBroker source)
{
//
var result =
!source.IsNull() &&
source != XEventBroker.None;
//
return result;
}
/// <summary>
/// Validate Event Service Configuration ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XEventServiceConfiguration source)
{
//
var result =
!source.IsNullOrDefault() &&
source.Broker.IsValid() &&
!source.Url.IsNullOrEmpty() &&
source.Url.IsValidUrl();
//
return result;
}
/// <summary>
/// Validate RabbitMQ Exchange Type ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XRabbitMQExchangeType source)
{
//
var result = source != XRabbitMQExchangeType.None;
//
return result;
}
/// <summary>
/// Validate RabbitMQ Configuration ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XRabbitMQOptions source)
{
//
var result =
!source.IsNullOrDefault() &&
!source.Exchange.IsNullOrEmpty() &&
source.ExchangeType.IsValid();
//
return result;
}
/// <summary>
/// Validate an Event Message ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XEventMessage source)
{
//
var result =
!source.IsNullOrDefault() &&
!source.Action.IsNullOrEmpty() &&
!source.Sender.IsNullOrEmpty();
//
return result;
}
/// <summary>
/// Validate an Event Request ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XEventRequest source)
{
//
var result =
!source.IsNullOrDefault() &&
!source.Topic.IsNullOrEmpty() &&
source.Message.IsValid();
//
return result;
}
/// <summary>
/// Validate ZeroMQ Pattern ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XZeroMQPattern source)
{
//
var result = source != XZeroMQPattern.None;
//
return result;
}
/// <summary>
/// Validate ZeroMQ Options ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsValid(this XZeroMQOptions source)
{
//
var result =
!source.IsNullOrDefault() &&
source.Pattern.IsValid();
//
return result;
}
/// <summary>
/// Validate an Array is Wildcard passed or not ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsAllowedWildcard(this string[] source)
{
//
var result =
source is not null &&
source.Length == 1 &&
source.Any(s =>
s == XEventServiceConstants.XWildcard);
//
return result;
}
/// <summary>
/// Validate a Collection is Valid for Kafka Consumer ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static bool IsAllowedCollection(
this string[] source
)
{
//
var result = source.HasChild();
//
return result;
}
/// <summary>
/// Validate Sender ...
/// </summary>
/// <param name="source"></param>
/// <param name="allowedSenders"></param>
/// <returns></returns>
public static bool IsValidSender(
this XEventMessage source,
string[] allowedSenders
)
{
//
// Validate Source and Collection ...
var result =
source.IsValid() &&
allowedSenders
.IsAllowedCollection();
if (!result)
{
return result;
}
//
// Validate Sender ...
result =
allowedSenders.IsAllowedWildcard() ||
allowedSenders.Any(x =>
x.Equals(
source.Sender,
System.StringComparison.InvariantCultureIgnoreCase));
//
return result;
}
/// <summary>
/// Validate Action ...
/// </summary>
/// <param name="source"></param>
/// <param name="allowedActions"></param>
/// <returns></returns>
public static bool IsValidAction(
this XEventMessage source,
string[] allowedActions
)
{
//
// Validate Source and Collection ...
var result =
source.IsValid() &&
allowedActions
.IsAllowedCollection();
if (!result)
{
return result;
}
//
// Validate Action ...
result =
allowedActions.IsAllowedWildcard() ||
allowedActions.Any(x =>
x.Equals(
source.Action,
System.StringComparison.InvariantCultureIgnoreCase));
//
return result;
}
/// <summary>
/// Completely Validate a Message, based on providing:
/// - Allowed Actions
/// - Allowed Senders
/// </summary>
/// <param name="source"></param>
/// <param name="allowedSenders"></param>
/// <param name="allowedActions"></param>
/// <returns></returns>
public static bool IsValid(
this XEventMessage source,
string[] allowedSenders,
string[] allowedActions
)
{
//
var result =
source.IsValid() &&
source.IsValidAction(allowedActions) &&
source.IsValidSender(allowedSenders);
//
return result;
}
/// <summary>
/// Parse Url and Extract Host address ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static string GetRabbitMQHost(this XEventServiceConfiguration source)
{
//
var result = string.Empty;
//
if (!source.IsValid())
{
XException.InvalidConfiguration.Throw();
}
//
var parts = source.Url.Split(':');
if (parts.Length != 2)
{
XException.InvalidConfiguration.Throw();
}
//
result = parts[0];
if (result.IsNullOrEmpty() ||
!result.IsValidUrl()
)
{
XException.InvalidConfiguration.Throw();
}
//
return result;
}
/// <summary>
/// Parse Url and Extract Port Number ...
/// </summary>
/// <param name="source"></param>
/// <returns></returns>
public static int GetRabbitMQPort(this XEventServiceConfiguration source)
{
//
int result = -1;
//
if (!source.IsValid())
{
XException.InvalidConfiguration.Throw();
}
//
var parts = source.Url.Split(':');
if (parts.Length != 2)
{
XException.InvalidConfiguration.Throw();
}
//
// Try Parse Port as Integer ...
int.TryParse(parts[1], out result);
if (result <= 0)
{
XException.InvalidConfiguration.Throw();
}
//
return result;
}
}
}
-11
View File
@@ -1,11 +0,0 @@
using xEventService.Models;
namespace xEventService.Interfaces {
public interface IXEventBus {
bool Publish<T> (T @event) where T : XEvent;
bool Subscribe<TEvent, THandler> ()
where TEvent : XEvent
where THandler : IXEventHandler<TEvent>;
}
}
-11
View File
@@ -1,11 +0,0 @@
using System.Threading.Tasks;
using xEventService.Models;
namespace xEventService.Interfaces {
public interface IXEventHandler<in TEvent> : IXEventHandler
where TEvent : XEvent {
Task HandleAsync (TEvent @event);
}
public interface IXEventHandler { }
}
+5
View File
@@ -0,0 +1,5 @@
namespace xEventService.Interfaces
{
public interface IXEventServiceProvider : IXProducerServiceBase
{ }
}
+8
View File
@@ -0,0 +1,8 @@
namespace xEventService.Interfaces
{
/// <summary>
/// a Service for Produce Kafka Messages ...
/// </summary>
public interface IXKafkaProducerService : IXProducerServiceBase
{ }
}
+35
View File
@@ -0,0 +1,35 @@
using System;
using System.Threading;
using xEventService.Models;
using System.Threading.Tasks;
using xEventService.Constants;
namespace xEventService.Interfaces
{
public interface IXProducerServiceBase : IDisposable
{
/// <summary>
/// Sending Message ...
/// </summary>
/// <param name="request">an Instance of <see cref="XEventRequest"/> to Provides Kafka Messaging requirement ...</param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
Task<bool> SendMessageAsync(
XEventRequest request,
CancellationToken cancellationToken = default
);
/// <summary>
/// Sending a Message ...
/// </summary>
/// <param name="message">an Instance of <see cref="XEventMessage"/> to Provides Kafka Messaging requirement ...</param>
/// <param name="topic">Topic for Sending Message ...</param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
Task<bool> SendMessageAsync(
XEventMessage message,
string topic = XEventServiceConstants.XDefaultTopic,
CancellationToken cancellationToken = default
);
}
}
+33
View File
@@ -0,0 +1,33 @@
using System.Threading;
using xEventService.Models;
using System.Threading.Tasks;
namespace xEventService.Interfaces
{
/// <summary>
/// RabbitMQ Producer Service Implementation ...
/// </summary>
public interface IXRabbitMQProducerService : IXProducerServiceBase
{
/// <summary>
/// Sending a Message to Default Exchange ...
/// </summary>
/// <param name="message">
/// an Instance of <see cref="XEventMessage"/> ...
/// </param>
/// <param name="exchange">
/// Target Exchange Name. if empty, uses default from configuration ...
/// </param>
/// <param name="routingKey">
/// Routing Key for message routing ...
/// </param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
Task<bool> SendMessageAsync(
XEventMessage message,
string exchange = null,
string routingKey = null,
CancellationToken cancellationToken = default
);
}
}
+5
View File
@@ -0,0 +1,5 @@
namespace xEventService.Interfaces
{
public interface IXZeroMQProducerService : IXProducerServiceBase
{ }
}
-224
View File
@@ -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 ...
/// <summary>
/// Close Channel ...
/// </summary>
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;
}
}
/// <summary>
/// Close Connection ...
/// </summary>
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;
}
}
/// <summary>
/// Publish a Message through Channel on Declare Queue ...
/// </summary>
/// <param name="message"></param>
/// <typeparam name="T"></typeparam>
/// <returns></returns>
public bool Publish<T> (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
);
}
/// <summary>
/// Dispose XChannel ...
/// </summary>
public void Dispose () {
//
this.CloseChannel ();
this.CloseConnection ();
}
#endregion
}
}
-5
View File
@@ -1,5 +0,0 @@
using xEventService.Base;
namespace xEventService.Models {
public class XEvent : XBaseEvent { }
}
+40
View File
@@ -0,0 +1,40 @@
using System;
using Newtonsoft.Json;
using System.Collections.Generic;
namespace xEventService.Models
{
/// <summary>
/// a Message Descriptor for Events ...
/// </summary>
public class XEventMessage
{
/// <summary>
/// Specified Action for Fire ...
/// </summary>
[JsonRequired]
public string Action { get; set; }
/// <summary>
/// Specified Action Sender ...
/// </summary>
[JsonRequired]
public string Sender { get; set; }
/// <summary>
/// Sending Time ...
/// </summary>
[JsonRequired]
public DateTime Offset { get; set; }
/// <summary>
/// Specified Message ...
/// </summary>
public string Message { get; set; }
/// <summary>
/// Metadata for Message ...
/// </summary>
public IDictionary<string, string> Payload { get; set; } = new Dictionary<string, string>();
}
}
+20
View File
@@ -0,0 +1,20 @@
using xEventService.Constants;
namespace xEventService.Models
{
/// <summary>
/// a Request for Sending an Event ...
/// </summary>
public class XEventRequest
{
/// <summary>
/// Specified Which Kafka Topic ...
/// </summary>
public string Topic { get; set; } = XEventServiceConstants.XDefaultTopic;
/// <summary>
/// Specified Kafka Meesage to Send ...
/// </summary>
public XEventMessage Message { get; set; }
}
}
+92
View File
@@ -0,0 +1,92 @@
using xEventService.Constants;
namespace xEventService.Models
{
/// <summary>
/// RabbitMQ-specific connection and messaging options ...
/// </summary>
public class XRabbitMQOptions
{
/// <summary>
/// Connection Username ...
/// </summary>
public string Username { get; set; } = "guest";
/// <summary>
/// Connection Password ...
/// </summary>
public string Password { get; set; } = "guest";
/// <summary>
/// Virtual Host ...
/// </summary>
public string VirtualHost { get; set; } = "/";
/// <summary>
/// Default Exchange Name ...
/// </summary>
public string Exchange { get; set; } = "x.saherelm.events";
/// <summary>
/// Default Exchange Type (fanout, direct, topic, headers) ...
/// </summary>
public XRabbitMQExchangeType ExchangeType { get; set; } = XRabbitMQExchangeType.Fanout;
/// <summary>
/// Default Queue Name ...
/// if empty, a dynamic queue will be generated ...
/// </summary>
public string Queue { get; set; } = string.Empty;
/// <summary>
/// Default Routing Key ...
/// </summary>
public string RoutingKey { get; set; } = string.Empty;
/// <summary>
/// Exchange Durable Flag ...
/// </summary>
public bool ExchangeDurable { get; set; } = true;
/// <summary>
/// Exchange AutoDelete Flag ...
/// </summary>
public bool ExchangeAutoDelete { get; set; } = false;
/// <summary>
/// Queue Durable Flag ...
/// </summary>
public bool QueueDurable { get; set; } = true;
/// <summary>
/// Queue Exclusive Flag ...
/// </summary>
public bool QueueExclusive { get; set; } = false;
/// <summary>
/// Queue AutoDelete Flag ...
/// </summary>
public bool QueueAutoDelete { get; set; } = false;
/// <summary>
/// Enable Automatic Recovery on Connection Loss ...
/// </summary>
public bool AutomaticRecoveryEnabled { get; set; } = true;
/// <summary>
/// Network Recovery Interval in Seconds ...
/// </summary>
public int NetworkRecoveryIntervalSeconds { get; set; } = 5;
/// <summary>
/// Prefetch Count for Consumer ...
/// </summary>
public ushort PrefetchCount { get; set; } = 1;
/// <summary>
/// Enable Manual Acknowledgement ...
/// </summary>
public bool ManualAck { get; set; } = true;
}
}
+28
View File
@@ -0,0 +1,28 @@
using xEventService.Constants;
namespace xEventService.Models
{
public class XZeroMQOptions
{
/// <summary>
/// Socket Pattern (PubSub, PushPull) ...
/// </summary>
public XZeroMQPattern Pattern { get; set; } = XZeroMQPattern.PubSub;
/// <summary>
/// Bind or Connect ...
/// True for Server (Bind), False for Client (Connect) ...
/// </summary>
public bool IsServer { get; set; } = false;
/// <summary>
/// High Water Mark for Sending ...
/// </summary>
public int SendHighWaterMark { get; set; } = 1000;
/// <summary>
/// High Water Mark for Receiving ...
/// </summary>
public int ReceiveHighWaterMark { get; set; } = 1000;
}
}
@@ -0,0 +1,65 @@
using System;
using xCommons.Providers;
using xCommons.Extensions;
using xExceptions.Constants;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
/// <summary>
/// Base Event Consumer Service ...
/// </summary>
public abstract class XBaseEventConsumerBackgroundService : XResumableBackgroundService
{
/// <summary>
/// Consumed Topic Messages ...
/// </summary>
protected readonly string topic;
/// <summary>
/// Allowed Senders ...
/// </summary>
protected readonly string[] allowedSenders;
/// <summary>
/// Allowed Actions ...
/// </summary>
protected readonly string[] allowedActions;
/// <summary>
/// Configuration ...
/// </summary>
protected readonly XEventServiceConfiguration configuration;
protected XBaseEventConsumerBackgroundService(
string[] allowedActions,
string[] allowedSenders,
IServiceProvider serviceProvider,
XEventServiceConfiguration configuration,
ILogger<XResumableBackgroundService> logger,
string topic = XEventServiceConstants.XDefaultTopic
) : base(serviceProvider, logger)
{
//
// Validate Actions ...
if (!allowedActions
.IsAllowedCollection() ||
!allowedSenders
.IsAllowedCollection() ||
!configuration.IsValid()
)
{
XException.InvalidArgs.Throw();
}
//
this.topic = topic;
this.configuration = configuration;
this.allowedActions = allowedActions;
this.allowedSenders = allowedSenders;
}
}
}
+31
View File
@@ -0,0 +1,31 @@
using xCommons.Extensions;
using xExceptions.Constants;
using xEventService.Extensions;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
public abstract class XBaseEventProducerService
{
protected readonly XEventServiceConfiguration configuration;
protected readonly ILogger<XBaseEventProducerService> logger;
protected XBaseEventProducerService(
XEventServiceConfiguration configuration,
ILogger<XBaseEventProducerService> logger
)
{
//
this.logger = logger;
this.configuration = configuration;
//
// Validate Configurations ...
if (!configuration.IsValid())
{
XException.InvalidConfiguration.Throw();
}
}
}
}
-260
View File
@@ -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<Type> events;
private readonly ILogger<XEventBus> logger;
private readonly IServiceScopeFactory serviceScopeFactory;
private readonly XEventServiceConfiguration configuration;
private readonly IDictionary<string, List<Type>> handlers;
#endregion
//
#region Constructors ...
public XEventBus (
ILoggerFactory loggerFactory,
IServiceScopeFactory serviceScopeFactory,
XEventServiceConfiguration configuration
) {
//
this.configuration = configuration;
this.serviceScopeFactory = serviceScopeFactory;
//
events = new List<Type> ();
handlers = new Dictionary<string, List<Type>> ();
logger = loggerFactory.CreateLogger<XEventBus> ();
}
#endregion
//
#region Abstracts ...
public bool Publish<T> (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<TEvent, THandler> ()
where TEvent : XEvent
where THandler : IXEventHandler<TEvent> {
//
// 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<Type> ());
}
//
// 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<TEvent> ();
//
return true;
}
#endregion
//
#region Private ...
/// <summary>
/// Consume Specific Event ...
/// </summary>
/// <typeparam name="TEvent"></typeparam>
private void ConsumeEvent<TEvent> ()
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);
}
/// <summary>
/// this Delegate method calls when a message recieved ...
/// </summary>
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
);
}
}
/// <summary>
/// Here we must call registered Event Handlers to Handle the Event ...
/// </summary>
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<XEvent>.HandleAsync))
.Invoke (handler, new object[] { @event }))
.ConfigureAwait (true);
}
}
}
#endregion
}
}
+159
View File
@@ -0,0 +1,159 @@
using System.Threading;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Interfaces;
using xEventService.Configuration;
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();
}
}
/// <summary>
/// Send Message to Broker Using XEventRequest instance ...
/// </summary>
/// <param name="request"></param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> 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;
}
/// <summary>
/// Send Message to Broker Using XMessage instance ...
/// </summary>
/// <param name="message"></param>
/// <param name="topic"></param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> 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;
}
/// <summary>
/// Dispose ...
/// </summary>
public void Dispose()
{
//
if (kafkaProducerService != null)
{
kafkaProducerService.Dispose();
}
//
if (zeroMQProducerService != null)
{
zeroMQProducerService.Dispose();
}
//
if (rabbitMQProducerService != null)
{
rabbitMQProducerService.Dispose();
}
}
}
}
+142
View File
@@ -0,0 +1,142 @@
using System;
using Confluent.Kafka;
using System.Threading;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
/// <summary>
/// an Abstract Kafka Consumer Service ...
/// </summary>
public abstract class XKafkaConsumerServiceBase : XBaseEventConsumerBackgroundService
{
/// <summary>
/// Kafka Consumer Class ...
/// </summary>
protected readonly IConsumer<Ignore, string> consumer;
protected XKafkaConsumerServiceBase(
string[] allowedActions,
string[] allowedSenders,
IServiceProvider serviceProvider,
XEventServiceConfiguration configuration,
ILogger<XKafkaConsumerServiceBase> logger,
string topic = XEventServiceConstants.XDefaultTopic
) : base(
topic: topic,
logger: logger,
configuration: configuration,
allowedActions: allowedActions,
allowedSenders: allowedSenders,
serviceProvider: serviceProvider
)
{
//
// Validating Configuration ...
if (configuration.Broker != XEventBroker.XKafka)
{
XException.InvalidConfiguration.Throw();
}
//
// Prepare Kafka Consumer Configuration ...
var config = new ConsumerConfig
{
EnableAutoCommit = false,
BootstrapServers = configuration.Url,
AutoOffsetReset = AutoOffsetReset.Earliest,
GroupId = XEventServiceConstants.XDefaultConsumerGroup,
};
consumer = new ConsumerBuilder<Ignore, string>(config).Build();
}
/// <summary>
/// Default Action Execution for Background Services ...
/// </summary>
/// <param name="cancellationToken"></param>
/// <returns></returns>
protected override async Task ExecuteAsync(
CancellationToken cancellationToken = default
)
{
//
// For Fix Blocking Synchronous ...
await Task.Yield();
using var cts = new CancellationTokenSource();
//
// Subscribe to Topic ...
consumer.Subscribe(topic);
logger.LogInformation($"Kafka Cosume Topic: {topic} ...");
//
// Consuming ...
while (!cancellationToken.IsCancellationRequested)
{
try
{
//
// Consuming Topic ...
var consumeResult = consumer.Consume(cancellationToken);
if (!consumeResult.IsNullOrDefault() &&
!consumeResult.Message.IsNullOrDefault()
)
{
//
var message = consumeResult
.Message
.Value
.FromJSON<XEventMessage>();
if (message.IsValid(
allowedSenders: allowedSenders,
allowedActions: allowedActions)
)
{
//
await DoWorkAsync(
payload: message,
cancellationToken: cancellationToken
);
}
}
}
catch (ConsumeException ex)
{
//
logger.LogError($"Kafka Consumer Result Parsing Error: {ex.Message} ...");
await Task.Delay(1000, cancellationToken);
}
catch (Exception ex)
{
//
logger.LogError($"Kafka Consumer Result Parsing Error: {ex.Message} ...");
await Task.Delay(1000, cancellationToken);
}
}
}
/// <summary>
/// Dispose Object ...
/// </summary>
public override void Dispose()
{
//
if (consumer != null)
{
consumer.Close();
consumer.Dispose();
}
//
GC.SuppressFinalize(this);
base.Dispose();
}
}
}
+140
View File
@@ -0,0 +1,140 @@
using System;
using Confluent.Kafka;
using System.Threading;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Interfaces;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
public class XKafkaProducerService : XBaseEventProducerService, IXKafkaProducerService
{
private readonly IProducer<Null, string> producer;
public XKafkaProducerService(
ILogger<XKafkaProducerService> logger,
XEventServiceConfiguration configuration
) : base(configuration, logger)
{
//
// Validate Configurations ...
if (!configuration.IsValid() ||
configuration.Broker != XEventBroker.XKafka
)
{
XException.InvalidConfiguration.Throw();
}
//
// Prepare Kafka Producer Configuration ...
var config = new ProducerConfig
{
Acks = Acks.All,
BootstrapServers = configuration.Url,
MessageSendMaxRetries = configuration.MaxRetries,
};
producer = new ProducerBuilder<Null, string>(config).Build();
}
/// <summary>
/// Sending Message ...
/// </summary>
/// <param name="request">an Instance of <see cref="XEventRequest"/> to Provides Kafka Messaging requirement ...</param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> SendMessageAsync(
XEventRequest request,
CancellationToken cancellationToken = default
)
{
//
// Validate ...
var result = request.IsValid();
if (!result)
{
XException.InvalidArgs.Throw();
}
//
result = await SendMessageAsync(
topic: request.Topic,
message: request.Message,
cancellationToken: cancellationToken
);
//
return result;
}
/// <summary>
/// Sending a Message ...
/// </summary>
/// <param name="message">an Instance of <see cref="XEventMessage"/> to Provides Kafka Messaging requirement ...</param>
/// <param name="topic">Topic for Sending Message ...</param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> SendMessageAsync(
XEventMessage message,
string topic = XEventServiceConstants.XDefaultTopic,
CancellationToken cancellationToken = default
)
{
//
// Validate ...
var result = message.IsValid() &&
!topic.IsNullOrEmpty();
if (!result)
{
XException.InvalidArgs.Throw();
}
//
try
{
//
var response = await producer.ProduceAsync(
topic: topic,
message: new Message<Null, string>
{
Value = message
.ToJSON(camelCase: true)
},
cancellationToken: cancellationToken
);
//
result = !response.IsNullOrDefault();
}
catch (Exception ex)
{
//
logger.LogError($"Kafka Produce Error: {ex.Message} ...");
result = false;
}
//
return result;
}
/// <summary>
/// Dispose Required Objects ...
/// </summary>
public void Dispose()
{
//
if (producer != null)
{
producer.Dispose();
}
//
GC.SuppressFinalize(this);
}
}
}
+277
View File
@@ -0,0 +1,277 @@
using System;
using System.Text;
using RabbitMQ.Client;
using System.Threading;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using System.Threading.Tasks;
using RabbitMQ.Client.Events;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
/// <summary>
/// an Abstract Kafka Consumer Service ...
/// </summary>
public abstract class XRabbitMQConsumerServiceBase : XBaseEventConsumerBackgroundService
{
//
protected IModel channel;
protected IConnection connection;
protected AsyncEventingBasicConsumer consumer;
protected XRabbitMQConsumerServiceBase(
string[] allowedActions,
string[] allowedSenders,
IServiceProvider serviceProvider,
XEventServiceConfiguration configuration,
ILogger<XRabbitMQConsumerServiceBase> logger,
string topic = XEventServiceConstants.XDefaultTopic
) : base(
topic: topic,
logger: logger,
configuration: configuration,
allowedActions: allowedActions,
allowedSenders: allowedSenders,
serviceProvider: serviceProvider
)
{
//
// Validating Configuration ...
if (!configuration.RabbitMQ.IsValid() ||
configuration.Broker != XEventBroker.XRabbitMQ
)
{
XException.InvalidConfiguration.Throw();
}
//
SetupConnection();
DeclareTopology(topic);
consumer = new AsyncEventingBasicConsumer(channel);
}
/// <summary>
/// Default Action Execution for Background Services ...
/// </summary>
/// <param name="cancellationToken"></param>
/// <returns></returns>
protected override async Task ExecuteAsync(
CancellationToken cancellationToken = default
)
{
//
// For Fix Blocking Synchronous ...
await Task.Yield();
using var cts = new CancellationTokenSource();
logger.LogInformation($"RabbitMQ Cosume Topic: {topic} ...");
//
// Subscribe to Topic ...
consumer.Received += async (model, ea) =>
{
//
try
{
//
var body = ea.Body.ToArray();
var json = Encoding.UTF8.GetString(body);
var message = json.FromJSON<XEventMessage>();
if (!message.IsNullOrDefault() &&
message.IsValid(
allowedSenders: allowedSenders,
allowedActions: allowedActions)
)
{
//
await DoWorkAsync(
payload: message,
cancellationToken: cancellationToken
);
}
//
if (configuration.RabbitMQ.ManualAck)
{
//
channel.BasicAck(
multiple: false,
deliveryTag: ea.DeliveryTag
);
}
}
catch (Exception ex)
{
//
logger.LogError($"RabbitMQ Cosume Error: {ex.Message} ...");
if (configuration.RabbitMQ.ManualAck)
{
//
channel.BasicNack(
requeue: true,
multiple: false,
deliveryTag: ea.DeliveryTag
);
}
}
};
//
consumer.Shutdown += (model, ea) =>
{
return Task.CompletedTask;
};
//
channel.BasicQos(
prefetchSize: 0,
global: false,
prefetchCount: configuration.RabbitMQ.PrefetchCount
);
//
channel.BasicConsume(
queue: GetQueueName(),
consumer: consumer,
autoAck: !configuration.RabbitMQ.ManualAck
);
//
while (!cancellationToken.IsCancellationRequested)
{
await Task.Delay(1000, cancellationToken);
}
}
/// <summary>
/// Overriding Stop Action ...
/// </summary>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public override async Task StopAsync(CancellationToken cancellationToken)
{
//
if (channel != null)
{
//
channel.Close();
channel.Dispose();
}
//
if (connection != null)
{
//
connection.Close();
connection.Dispose();
}
//
await base.StopAsync(cancellationToken);
}
/// <summary>
/// Dispose Implementation ...
/// </summary>
public override void Dispose()
{
//
if (channel != null)
{
channel.Dispose();
}
//
if (connection != null)
{
connection.Dispose();
}
//
GC.SuppressFinalize(this);
base.Dispose();
}
//
#region Private ...
private string GetQueueName()
{
//
var result = string.Empty;
if (!configuration.RabbitMQ.Queue.IsNullOrEmpty())
{
result = configuration.RabbitMQ.Queue;
}
//
if (result.IsNullOrEmpty())
{
//
var exchangeName = topic ?? configuration.RabbitMQ.Exchange;
var digits = Guid.NewGuid().ToString().GetDigits();
result = $"{exchangeName}[{digits}]";
}
//
return result;
}
private void SetupConnection()
{
var factory = new ConnectionFactory
{
DispatchConsumersAsync = true,
Port = configuration.GetRabbitMQPort(),
HostName = configuration.GetRabbitMQHost(),
UserName = configuration.RabbitMQ.Username,
Password = configuration.RabbitMQ.Password,
NetworkRecoveryInterval = TimeSpan.FromSeconds(
configuration.RabbitMQ.NetworkRecoveryIntervalSeconds
),
VirtualHost = configuration.RabbitMQ.VirtualHost,
AutomaticRecoveryEnabled = configuration.RabbitMQ.AutomaticRecoveryEnabled,
};
//
connection = factory.CreateConnection();
channel = connection.CreateModel();
}
private void DeclareTopology(string topic)
{
//
channel.ExchangeDeclare(
arguments: null,
exchange: topic ?? configuration.RabbitMQ.Exchange,
durable: configuration.RabbitMQ.ExchangeDurable,
autoDelete: configuration.RabbitMQ.ExchangeAutoDelete,
type: configuration.RabbitMQ.ExchangeType.GetStringValue()
);
//
var queueName = GetQueueName();
//
channel.QueueDeclare(
arguments: null,
queue: queueName,
durable: configuration.RabbitMQ.QueueDurable,
exclusive: configuration.RabbitMQ.QueueExclusive,
autoDelete: configuration.RabbitMQ.QueueAutoDelete
);
//
channel.QueueBind(
arguments: null,
queue: queueName,
exchange: topic ?? configuration.RabbitMQ.Exchange,
routingKey: configuration.RabbitMQ.RoutingKey
);
}
#endregion
}
}
+315
View File
@@ -0,0 +1,315 @@
using System;
using System.Text;
using RabbitMQ.Client;
using System.Threading;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Interfaces;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
/// <summary>
/// RabbitMQ Producer Service Implementation ...
/// </summary>
public class XRabbitMQProducerService : XBaseEventProducerService, IXRabbitMQProducerService
{
private IModel channel;
private IConnection connection;
private readonly object connectionLock = new();
public XRabbitMQProducerService(
ILogger<XRabbitMQProducerService> logger,
XEventServiceConfiguration configuration
) : base(configuration, logger)
{
//
// Validate Configurations ...
if (!configuration.IsValid() ||
!configuration.RabbitMQ.IsValid() ||
configuration.Broker != XEventBroker.XRabbitMQ
)
{
XException.InvalidConfiguration.Throw();
}
//
// Initialize Connection ...
CreateConnection();
}
/// <summary>
/// Sending Message with Full Request Descriptor ...
/// </summary>
/// <param name="request"></param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> SendMessageAsync(
XEventRequest request,
CancellationToken cancellationToken = default
)
{
//
// Validate ...
var result = request.IsValid();
if (!result)
{
XException.InvalidArgs.Throw();
}
//
result = await SendMessageAsync(
message: request.Message,
exchange: request.Topic,
cancellationToken: cancellationToken,
routingKey: configuration.RabbitMQ.RoutingKey
);
//
return result;
}
/// <summary>
/// Sending a Message to Topic ...
/// </summary>
/// <param name="message"></param>
/// <param name="topic"></param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> SendMessageAsync(
XEventMessage message,
string topic = XEventServiceConstants.XDefaultTopic,
CancellationToken cancellationToken = default
)
{
//
var result = await SendMessageAsync(
request: new XEventRequest
{
Topic = topic,
Message = message
},
cancellationToken: cancellationToken
);
//
return result;
}
/// <summary>
/// Sending a Message to Exchange ...
/// </summary>
/// <param name="message"></param>
/// <param name="exchange"></param>
/// <param name="routingKey"></param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> SendMessageAsync(
XEventMessage message,
string exchange = null,
string routingKey = null,
CancellationToken cancellationToken = default
)
{
//
// Validate ...
var result = message.IsValid();
if (!result)
{
XException.InvalidArgs.Throw();
}
//
// Normalize Exchange ...
if (exchange.IsNullOrEmpty())
{
exchange = configuration.RabbitMQ.Exchange;
}
//
// Normalize Routing Key ...
if (routingKey.IsNull())
{
routingKey = configuration.RabbitMQ.RoutingKey;
}
//
try
{
//
// Ensure Connection is Alive ...
EnsureChannel();
//
// Serialize Message ...
var body = Encoding.UTF8.GetBytes(
message.ToJSON(camelCase: true)
);
//
// Prepare Message Properties ...
var properties = channel.CreateBasicProperties();
properties.Persistent = true;
properties.ContentType = "application/json";
properties.Timestamp = new AmqpTimestamp(
DateTimeOffset.UtcNow.ToUnixTimeSeconds()
);
properties.MessageId = Guid.NewGuid().ToString();
properties.Type = message.Action;
properties.AppId = message.Sender;
//
// Publish Message ...
await Task.Run(() =>
{
//
channel.BasicPublish(
exchange: exchange,
routingKey: routingKey,
mandatory: false,
basicProperties: properties,
body: body
);
},
cancellationToken: cancellationToken);
//
result = true;
}
catch (Exception ex)
{
//
logger.LogError($"RabbitMQ Produce Error: {ex.Message} ...");
result = false;
}
//
return result;
}
/// <summary>
/// Dispose Object ...
/// </summary>
public void Dispose()
{
//
try
{
//
if (channel != null)
{
//
channel.Close();
channel.Dispose();
channel = null;
}
//
if (connection != null)
{
//
connection.Close();
connection.Dispose();
connection = null;
}
}
catch
{ }
//
GC.SuppressFinalize(this);
}
//
#region Private ...
/// <summary>
/// Create Connection to RabbitMQ Server ...
/// </summary>
private void CreateConnection()
{
//
var factory = new ConnectionFactory
{
Port = configuration.GetRabbitMQPort(),
HostName = configuration.GetRabbitMQHost(),
UserName = configuration.RabbitMQ.Username,
Password = configuration.RabbitMQ.Password,
VirtualHost = configuration.RabbitMQ.VirtualHost,
AutomaticRecoveryEnabled = configuration.RabbitMQ.AutomaticRecoveryEnabled,
NetworkRecoveryInterval = TimeSpan.FromSeconds(
configuration.RabbitMQ.NetworkRecoveryIntervalSeconds
),
DispatchConsumersAsync = true
};
//
connection = factory.CreateConnection();
}
/// <summary>
/// Check Connection is Alive ...
/// </summary>
private bool IsConnected()
{
//
var result =
!connection.IsNullOrDefault() &&
connection.IsOpen &&
!channel.IsNull() &&
channel.IsOpen;
//
return result;
}
/// <summary>
/// Ensure Channel is Open and Exchange/Queue Declared ...
/// </summary>
private void EnsureChannel()
{
//
if (IsConnected())
{
return;
}
//
lock (connectionLock)
{
//
// Double Check after Lock ...
if (IsConnected())
{
return;
}
//
// Recreate Connection if Dead ...
if (connection is null || !connection.IsOpen)
{
CreateConnection();
}
//
// Create Channel ...
channel = connection.CreateModel();
//
// Declare Exchange ...
channel.ExchangeDeclare(
arguments: null,
exchange: configuration.RabbitMQ.Exchange,
durable: configuration.RabbitMQ.ExchangeDurable,
autoDelete: configuration.RabbitMQ.ExchangeAutoDelete,
type: configuration.RabbitMQ.ExchangeType.GetStringValue()
);
}
}
#endregion
}
}
+206
View File
@@ -0,0 +1,206 @@
using NetMQ;
using System;
using System.Text;
using NetMQ.Sockets;
using System.Threading;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
/// <summary>
/// an Abstract ZeroMQ Consumer Service ...
/// </summary>
public abstract class XZeroMQConsumerServiceBase : XBaseEventConsumerBackgroundService
{
//
private NetMQSocket socket;
protected XZeroMQConsumerServiceBase(
string[] allowedActions,
string[] allowedSenders,
IServiceProvider serviceProvider,
XEventServiceConfiguration configuration,
ILogger<XZeroMQConsumerServiceBase> logger,
string topic = XEventServiceConstants.XDefaultTopic
) : base(
topic: topic,
logger: logger,
configuration: configuration,
allowedActions: allowedActions,
allowedSenders: allowedSenders,
serviceProvider: serviceProvider
)
{
//
// Validating Configuration ...
if (!configuration.ZeroMQ.IsValid() ||
configuration.Broker != XEventBroker.XZeroMQ)
{
XException.InvalidConfiguration.Throw();
}
//
// Initialize ...
InitializeSocket();
}
/// <summary>
/// Default Action Execution for Background Services ...
/// </summary>
/// <param name="cancellationToken"></param>
/// <returns></returns>
protected override async Task ExecuteAsync(
CancellationToken cancellationToken = default
)
{
//
await Task.Yield();
logger.LogInformation($"ZeroMQ Cosume Topic: {topic} ...");
while (!cancellationToken.IsCancellationRequested)
{
//
try
{
//
string receivedTopic = string.Empty;
byte[] body = null;
//
if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
{
//
if (socket.TryReceiveFrameString(TimeSpan.FromMilliseconds(100), out receivedTopic))
{
//
if (socket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(100), out body))
{
// Process
}
else
{
logger.LogError($"ZeroMQ Cosume Error ...");
}
}
}
else
{
if (socket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(100), out body))
{
// Process
}
else
{
logger.LogError($"ZeroMQ Cosume Error ...");
}
}
//
if (body != null && body.Length > 0)
{
//
var json = Encoding.UTF8.GetString(body);
var message = json.FromJSON<XEventMessage>();
if (!message.IsNullOrDefault() &&
message.IsValid(
allowedSenders: allowedSenders,
allowedActions: allowedActions)
)
{
//
await DoWorkAsync(
payload: message,
cancellationToken: cancellationToken
);
}
}
else
{
await Task.Delay(10, cancellationToken);
}
}
catch (Exception ex)
{
//
logger.LogError($"ZeroMQ Cosume Error: {ex.Message} ...");
await Task.Delay(1000, cancellationToken);
}
}
}
public override async Task StopAsync(CancellationToken cancellationToken)
{
//
if (socket != null)
{
socket.Close();
}
//
await base.StopAsync(cancellationToken);
}
/// <summary>
/// Dispose Object ...
/// </summary>
public override void Dispose()
{
//
if (socket != null)
{
socket.Dispose();
}
//
GC.SuppressFinalize(this);
base.Dispose();
}
//
#region Private ...
private void InitializeSocket()
{
//
if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
{
//
socket = new SubscriberSocket();
socket.Options.ReceiveHighWatermark = configuration.ZeroMQ.ReceiveHighWaterMark;
//
// Subscribe to specific topic or all
if (topic.IsNullOrEmpty() || topic == XEventServiceConstants.XWildcard)
{
((SubscriberSocket)socket).SubscribeToAnyTopic();
}
else
{
((SubscriberSocket)socket).Subscribe(topic);
}
}
else
{
//
socket = new PullSocket();
socket.Options.ReceiveHighWatermark = configuration.ZeroMQ.ReceiveHighWaterMark;
}
//
if (configuration.ZeroMQ.IsServer)
{
socket.Bind(configuration.Url);
}
else
{
socket.Connect(configuration.Url);
}
}
#endregion
}
}
+198
View File
@@ -0,0 +1,198 @@
using NetMQ;
using System;
using System.Text;
using NetMQ.Sockets;
using System.Threading;
using xCommons.Extensions;
using xEventService.Models;
using xExceptions.Constants;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Interfaces;
using xEventService.Configuration;
using Microsoft.Extensions.Logging;
namespace xEventService.Providers
{
public class XZeroMQProducerService : XBaseEventProducerService, IXZeroMQProducerService
{
//
private NetMQSocket socket;
private readonly object socketLock = new();
public XZeroMQProducerService(
XEventServiceConfiguration configuration,
ILogger<XZeroMQProducerService> logger
) : base(configuration, logger)
{
//
// Validate Configurations ...
if (!configuration.IsValid() ||
!configuration.ZeroMQ.IsValid() ||
configuration.Broker != XEventBroker.XZeroMQ
)
{
XException.InvalidConfiguration.Throw();
}
//
InitializeSocket();
}
/// <summary>
/// Sending Message with Full Request Descriptor ...
/// </summary>
/// <param name="request"></param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> SendMessageAsync(
XEventRequest request,
CancellationToken cancellationToken = default
)
{
//
var result = request.IsValid();
if (!result)
{
XException.InvalidArgs.Throw();
}
//
result = await SendMessageAsync(
message: request.Message,
topic: request.Topic,
cancellationToken: cancellationToken
);
//
return result;
}
/// <summary>
/// Sending a Message to Topic ...
/// </summary>
/// <param name="message"></param>
/// <param name="topic"></param>
/// <param name="cancellationToken"></param>
/// <returns></returns>
public async Task<bool> SendMessageAsync(
XEventMessage message,
string topic = XEventServiceConstants.XDefaultTopic,
CancellationToken cancellationToken = default
)
{
//
var result = message.IsValid();
if (!result)
{
XException.InvalidArgs.Throw();
}
//
// Normalize Topic ...
if (topic.IsNullOrEmpty())
{
topic = XEventServiceConstants.XDefaultTopic;
}
//
try
{
//
var json = message.ToJSON(camelCase: true);
var body = Encoding.UTF8.GetBytes(json);
await Task.Run(() =>
{
//
lock (socketLock)
{
//
if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
{
socket.SendMoreFrame(topic).SendFrame(body);
}
else
{
socket.SendFrame(body);
}
}
}, cancellationToken);
//
result = true;
}
catch (Exception)
{
result = false;
}
//
return result;
}
/// <summary>
/// Dispose Object ...
/// </summary>
public void Dispose()
{
//
try
{
//
lock (socketLock)
{
//
if (socket != null)
{
//
socket.Close();
socket.Dispose();
socket = null;
}
}
}
catch { }
//
GC.SuppressFinalize(this);
}
//
#region Private ...
/// <summary>
/// Initialize Socket ...
/// </summary>
private void InitializeSocket()
{
//
lock (socketLock)
{
//
if (configuration.ZeroMQ.Pattern == XZeroMQPattern.PubSub)
{
//
socket = new PublisherSocket();
socket.Options.SendHighWatermark = configuration.ZeroMQ.SendHighWaterMark;
}
else
{
//
socket = new PushSocket();
socket.Options.SendHighWatermark = configuration.ZeroMQ.SendHighWaterMark;
}
//
if (configuration.ZeroMQ.IsServer)
{
socket.Bind(configuration.Url);
}
else
{
socket.Connect(configuration.Url);
}
}
}
#endregion
}
}
+1 -1
View File
@@ -3,6 +3,6 @@
<packageSources>
<!-- <add key="megan" value="https://hub.megan.ir/nuget/index.json" /> -->
<add key="nuget" value="https://api.nuget.org/v3/index.json" protocolVersion="3" />
<add key="baget" value="https://nuget.saherelmhub.ir/v3/index.json" protocolVersion="3" disableTLSCertificateValidation="true" />
<!-- <add key="baget" value="https://nuget.saherelmhub.ir/v3/index.json" protocolVersion="3" disableTLSCertificateValidation="true" /> -->
</packageSources>
</configuration>
+4 -2
View File
@@ -35,14 +35,16 @@
<!-- Dependencies -->
<ItemGroup>
<PackageReference Include="NetMQ" Version="4.0.4.3" />
<PackageReference Include="RabbitMQ.Client" Version="6.2.2" />
<PackageReference Include="Confluent.Kafka" Version="2.14.2" />
</ItemGroup>
<!-- For XML Documentation Support -->
<PropertyGroup>
<CopyLocalLockFileAssemblies>true</CopyLocalLockFileAssemblies>
<GenerateDocumentationFile>true</GenerateDocumentationFile>
<NoWarn>$(NoWarn);1591</NoWarn>
<GenerateDocumentationFile>true</GenerateDocumentationFile>
<CopyLocalLockFileAssemblies>true</CopyLocalLockFileAssemblies>
</PropertyGroup>
</Project>