Files
xSaherelmWorkspace/Documents/GeneratedContext/xDashboard-xEventService-20260819_143307.txt
T
2026-08-19 15:37:05 +03:30

2212 lines
64 KiB
Plaintext

### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Configuration\XEventServiceConfiguration.cs
using xEventService.Constants;
using xEventService.Models;
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
{
/// <summary>
/// Message Broker Address ...
/// </summary>
public string Url { get; set; }
/// <summary>
/// Max Retries to Sending Message ...
/// </summary>
public int MaxRetries { get; set; } = 3;
/// <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();
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\ConfigurationNodeNames.cs
namespace xEventService.Constants
{
public partial struct ConfigurationNodeNames {
public const string EVENT_SERVICE_NODE_NAME = "EventServiceConfiguration";
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\XEventBroker.cs
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,
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\XEventServiceConstants.cs
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 = "*";
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\XRabbitMQExchangeType.cs
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,
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Constants\XZeroMQPattern.cs
using xExceptions.Attributes;
namespace xEventService.Constants
{
public enum XZeroMQPattern
{
[StringValue(XEventServiceConstants.None)]
None,
[StringValue(XEventServiceConstants.PubSub)]
PubSub,
[StringValue(XEventServiceConstants.PushPull)]
PushPull
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\DI\XDIHelperExtension.cs
using System;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
using xEventService.Extensions;
using xExceptions.Constants;
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
)
{
//
var configSection = source
.GetSection(ConfigurationNodeNames.EVENT_SERVICE_NODE_NAME);
var result = configSection.Get<XEventServiceConfiguration>();
//
return result;
}
/// <summary>
/// Register Event Service Configurations ...
/// </summary>
/// <param name="services"></param>
/// <param name="configuration"></param>
public static void AddXEventServiceConfigurations(
this IServiceCollection services,
IConfiguration configuration
)
{
//
var config = configuration.GetXEventServiceConfigurations();
services.AddXEventServiceConfigurations(config);
}
/// <summary>
/// Register Event Service Configurations ...
/// </summary>
/// <param name="services"></param>
/// <param name="configuration"></param>
public static void AddXEventServiceConfigurations(
this IServiceCollection services,
XEventServiceConfiguration configuration
)
{
//
// Validate ...
if (!configuration.IsValid())
{
//
Log("Registration Failed, due Invalid Configuration ...");
XException.InvalidConfiguration.Throw();
}
//
// Register Configuration as Singleton ...
services.AddSingleton(configuration);
}
//
#region Private ...
/// <summary>
/// a LogTag for Service ...
/// </summary>
private static string XLogTag = XEventServiceConstants.XEventServiceDILogTag;
/// <summary>
/// print a log in Console ...
/// </summary>
/// <param name="message"></param>
private static void Log(string message)
{
Console.WriteLine($"{XLogTag} => {message}");
}
#endregion
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Extensions\XEventServiceExtensions.cs
using System.Linq;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
using xEventService.Models;
using xExceptions.Constants;
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;
}
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXKafkaProducerService.cs
namespace xEventService.Interfaces
{
/// <summary>
/// a Service for Produce Kafka Messages ...
/// </summary>
public interface IXKafkaProducerService : IXProducerServiceBase
{ }
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXProducerServiceBase.cs
using System;
using System.Threading;
using System.Threading.Tasks;
using xEventService.Constants;
using xEventService.Models;
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
);
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXRabbitMQProducerService.cs
using System;
using System.Threading;
using System.Threading.Tasks;
using xEventService.Models;
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
);
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Interfaces\IXZeroMQProducerService.cs
namespace xEventService.Interfaces
{
public interface IXZeroMQProducerService : IXProducerServiceBase
{ }
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Models\XEventMessage.cs
using System;
using System.Collections.Generic;
using Newtonsoft.Json;
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>();
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Models\XEventRequest.cs
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; }
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Models\XRabbitMQOptions.cs
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;
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Models\XZeroMQOptions.cs
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;
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XBaseEventConsumerBackgroundService.cs
using System;
using Microsoft.Extensions.Logging;
using xCommons.Extensions;
using xCommons.Providers;
using xEventService.Configuration;
using xEventService.Constants;
using xEventService.Extensions;
using xExceptions.Constants;
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>
public readonly string[] allowedSenders;
/// <summary>
/// Allowed Actions ...
/// </summary>
public 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(
logger: logger,
serviceProvider: serviceProvider
)
{
//
// Validate Actions ...
if (!allowedActions
.IsAllowedCollection() ||
!allowedSenders
.IsAllowedCollection() ||
!configuration.IsValid()
)
{
XException.InvalidArgs.Throw();
}
//
this.topic = topic;
this.configuration = configuration;
this.allowedActions = allowedActions;
this.allowedSenders = allowedSenders;
}
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XBaseEventProducerService.cs
using Microsoft.Extensions.Logging;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Extensions;
using xExceptions.Constants;
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();
}
}
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XKafkaConsumerServiceBase.cs
using System;
using System.Threading;
using System.Threading.Tasks;
using Confluent.Kafka;
using Microsoft.Extensions.Logging;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
using xExceptions.Constants;
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);
//
// Consuming ...
while (!cancellationToken.IsCancellationRequested)
{
try
{
//
// Consuming Topic ...
var consumeResult = consumer.Consume(cancellationToken);
if (!consumeResult.IsNullOrDefault() &&
!consumeResult.Message.IsNullOrDefault()
)
{
//
await DoWorkAsync(
payload: consumeResult.Message,
cancellationToken: cancellationToken
);
}
}
catch (ConsumeException)
{
await Task.Delay(1000, cancellationToken);
}
catch (Exception)
{
await Task.Delay(1000, cancellationToken);
}
}
}
/// <summary>
/// Dispose Object ...
/// </summary>
public override void Dispose()
{
//
consumer.Close();
consumer.Dispose();
GC.SuppressFinalize(this);
base.Dispose();
}
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XKafkaProducerService.cs
using xCommons.Extensions;
using xExceptions.Constants;
using xEventService.Extensions;
using xEventService.Interfaces;
using xEventService.Configuration;
using xEventService.Models;
using Confluent.Kafka;
using Microsoft.Extensions.Logging;
using System.Threading.Tasks;
using System.Threading;
using xEventService.Constants;
using System;
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(
message: request.Message,
topic: request.Topic,
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
{
result = false;
}
//
return result;
}
/// <summary>
/// Dispose Required Objects ...
/// </summary>
public void Dispose()
{
//
producer.Dispose();
GC.SuppressFinalize(this);
}
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XRabbitMQConsumerServiceBase.cs
using System;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Models;
using xExceptions.Constants;
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();
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();
//
// 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())
{
//
await DoWorkAsync(
payload: message,
cancellationToken: cancellationToken
);
}
//
if (configuration.RabbitMQ.ManualAck)
{
//
channel.BasicAck(
multiple: false,
deliveryTag: ea.DeliveryTag
);
}
}
catch
{
//
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)
{
//
channel.Close();
channel.Dispose();
connection.Close();
connection.Dispose();
//
await base.StopAsync(cancellationToken);
}
/// <summary>
/// Dispose Implementation ...
/// </summary>
public override void Dispose()
{
//
channel.Dispose();
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 digits = Guid.NewGuid().ToString().GetDigits();
result = $"{configuration.RabbitMQ.Exchange}[{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()
{
//
channel.ExchangeDeclare(
arguments: null,
exchange: 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: configuration.RabbitMQ.Exchange,
routingKey: configuration.RabbitMQ.RoutingKey
);
}
#endregion
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XRabbitMQProducerService.cs
using System;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Interfaces;
using xEventService.Models;
using xExceptions.Constants;
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 ...
await EnsureChannelAsync(cancellationToken);
//
// 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
{
result = false;
}
//
return result;
}
/// <summary>
/// Dispose Object ...
/// </summary>
public void Dispose()
{
//
try
{
//
channel.Close();
channel.Dispose();
channel = 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 Connection Established ...
/// </summary>
private async Task<bool> EnsureConnectionAsync()
{
//
try
{
await EnsureChannelAsync();
return true;
}
catch
{
return false;
}
}
/// <summary>
/// Ensure Channel is Open and Exchange/Queue Declared ...
/// </summary>
private async Task EnsureChannelAsync(
CancellationToken cancellationToken = default
)
{
//
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
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XZeroMQConsumerServiceBase.cs
using System;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using NetMQ;
using NetMQ.Sockets;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
using xEventService.Models;
using xExceptions.Constants;
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.Broker != XEventBroker.XKafka)
{
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();
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
{
if (socket.TryReceiveFrameBytes(TimeSpan.FromMilliseconds(100), out body))
{
// Process
}
}
if (body != null && body.Length > 0)
{
//
var json = Encoding.UTF8.GetString(body);
var message = json.FromJSON<XEventMessage>();
if (!message.IsNullOrDefault())
{
//
await DoWorkAsync(
payload: message,
cancellationToken: cancellationToken
);
}
}
else
{
await Task.Delay(10, cancellationToken);
}
}
catch (Exception)
{
await Task.Delay(1000, cancellationToken);
}
}
}
public override async Task StopAsync(CancellationToken cancellationToken)
{
//
socket.Close();
await base.StopAsync(cancellationToken);
}
/// <summary>
/// Dispose Object ...
/// </summary>
public override void Dispose()
{
//
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
}
}
### FILE: C:\Users\hkhazaeeasl\Documents\Projets\XProjects\xSaherelmWorkspace\Modules\xEventService\Providers\XZeroMQProducerService.cs
using System;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using NetMQ;
using NetMQ.Sockets;
using xCommons.Extensions;
using xEventService.Configuration;
using xEventService.Constants;
using xEventService.Extensions;
using xEventService.Interfaces;
using xEventService.Models;
using xExceptions.Constants;
namespace xEventService.Providers
{
public class 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
}
}