diff --git a/Configuration/XEventServiceConfiguration.cs b/Configuration/XEventServiceConfiguration.cs index 157dbac..8e5df4c 100644 --- a/Configuration/XEventServiceConfiguration.cs +++ b/Configuration/XEventServiceConfiguration.cs @@ -1,6 +1,7 @@ using System.Collections.Generic; using RabbitMQ.Client; using xEventService.Constants; +using xEventService.Models; namespace xEventService.Configuration { @@ -24,5 +25,11 @@ namespace xEventService.Configuration /// Allowed Message Broker ... /// public XEventBroker Broker { get; set; } = XEventBroker.None; + + /// + /// RabbitMQ Specific Options ... + /// Only used when Broker is XRabbitMQ ... + /// + public XRabbitMQOptions RabbitMQ { get; set; } = new XRabbitMQOptions(); } } \ No newline at end of file diff --git a/Constants/XEventServiceConstants.cs b/Constants/XEventServiceConstants.cs index 538258d..e46022c 100644 --- a/Constants/XEventServiceConstants.cs +++ b/Constants/XEventServiceConstants.cs @@ -9,6 +9,9 @@ namespace xEventService.Constants // // Defaults ... public const string XDefaultTopic = "XSaherelm"; - public const string XDefaultConsumerGroup = "XSaherElmGroup"; + public const string XDefaultConsumerGroup = "XSaherElmGroup"; + + // + public const string XWildcard = "*"; } } \ No newline at end of file diff --git a/Constants/XRabbitMQExchangeType.cs b/Constants/XRabbitMQExchangeType.cs new file mode 100644 index 0000000..621ee9b --- /dev/null +++ b/Constants/XRabbitMQExchangeType.cs @@ -0,0 +1,37 @@ +using xExceptions.Attributes; + +namespace xEventService.Constants +{ + /// + /// Available RabbitMQ Exchange Types ... + /// + 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"; + } + + /// + /// RabbitMQ Exchange Type Enum ... + /// + public enum XRabbitMQExchangeType + { + [StringValue(XRabbitMQExchangeTypes.None)] + None, + + [StringValue(XRabbitMQExchangeTypes.Direct)] + Direct, + + [StringValue(XRabbitMQExchangeTypes.Fanout)] + Fanout, + + [StringValue(XRabbitMQExchangeTypes.Topic)] + Topic, + + [StringValue(XRabbitMQExchangeTypes.Headers)] + Headers, + } +} \ No newline at end of file diff --git a/Extensions/XEventServiceExtensions.cs b/Extensions/XEventServiceExtensions.cs index 421a7c8..f6c1c29 100644 --- a/Extensions/XEventServiceExtensions.cs +++ b/Extensions/XEventServiceExtensions.cs @@ -1,7 +1,9 @@ +using System.Linq; using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; using xEventService.Models; +using xExceptions.Constants; namespace xEventService.Extensions { @@ -44,6 +46,37 @@ namespace xEventService.Extensions return result; } + /// + /// Validate RabbitMQ Exchange Type ... + /// + /// + /// + public static bool IsValid(this XRabbitMQExchangeType source) + { + // + var result = source != XRabbitMQExchangeType.None; + + // + return result; + } + + /// + /// Validate RabbitMQ Configuration ... + /// + /// + /// + public static bool IsValid(this XRabbitMQOptions source) + { + // + var result = + !source.IsNullOrDefault() && + !source.Exchange.IsNullOrEmpty() && + !source.ExchangeType.IsValid(); + + // + return result; + } + /// /// Validate an Event Message ... /// @@ -77,5 +110,110 @@ namespace xEventService.Extensions // return result; } + + /// + /// Validate an Array is Wildcard passed or not ... + /// + /// + /// + public static bool IsAllowedWildcard(this string[] source) + { + // + var result = + source is not null && + source.Length == 1 && + source.Any(s => + s == XEventServiceConstants.XWildcard); + + // + return result; + } + + /// + /// Validate a Collection is Valid for Kafka Consumer ... + /// + /// + /// + public static bool IsAllowedCollection( + this string[] source + ) + { + // + var result = source.HasChild(); + + // + return result; + } + + /// + /// Parse Url and Extract Host address ... + /// + /// + /// + 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; + } + + /// + /// Parse Url and Extract Port Number ... + /// + /// + /// + 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; + } } } \ No newline at end of file diff --git a/Interfaces/IXRabbitMQProducerService.cs b/Interfaces/IXRabbitMQProducerService.cs new file mode 100644 index 0000000..2bd0806 --- /dev/null +++ b/Interfaces/IXRabbitMQProducerService.cs @@ -0,0 +1,48 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using xEventService.Models; + +namespace xEventService.Interfaces +{ + /// + /// RabbitMQ Producer Service Implementation ... + /// + public interface IXRabbitMQProducerService : IDisposable + { + /// + /// Sending Message with Full Request Descriptor ... + /// + /// + /// an Instance of . + /// Topic is used as Exchange Name ... + /// + /// + /// + Task SendMessageAsync( + XEventRequest request, + CancellationToken cancellationToken = default + ); + + /// + /// Sending a Message to Default Exchange ... + /// + /// + /// an Instance of ... + /// + /// + /// Target Exchange Name. if empty, uses default from configuration ... + /// + /// + /// Routing Key for message routing ... + /// + /// + /// + Task SendMessageAsync( + XEventMessage message, + string exchange = null, + string routingKey = null, + CancellationToken cancellationToken = default + ); + } +} \ No newline at end of file diff --git a/Models/XRabbitMQOptions.cs b/Models/XRabbitMQOptions.cs new file mode 100644 index 0000000..1b4584d --- /dev/null +++ b/Models/XRabbitMQOptions.cs @@ -0,0 +1,93 @@ +using RabbitMQ.Client; +using xEventService.Constants; + +namespace xEventService.Models +{ + /// + /// RabbitMQ-specific connection and messaging options ... + /// + public class XRabbitMQOptions + { + + /// + /// Connection Username ... + /// + public string Username { get; set; } = "guest"; + + /// + /// Connection Password ... + /// + public string Password { get; set; } = "guest"; + + /// + /// Virtual Host ... + /// + public string VirtualHost { get; set; } = "/"; + + /// + /// Default Exchange Name ... + /// + public string Exchange { get; set; } = "x.saherelm.events"; + + /// + /// Default Exchange Type (fanout, direct, topic, headers) ... + /// + public XRabbitMQExchangeType ExchangeType { get; set; } = XRabbitMQExchangeType.Fanout; + + /// + /// Default Queue Name ... + /// if empty, a dynamic queue will be generated ... + /// + public string Queue { get; set; } = string.Empty; + + /// + /// Default Routing Key ... + /// + public string RoutingKey { get; set; } = string.Empty; + + /// + /// Exchange Durable Flag ... + /// + public bool ExchangeDurable { get; set; } = true; + + /// + /// Exchange AutoDelete Flag ... + /// + public bool ExchangeAutoDelete { get; set; } = false; + + /// + /// Queue Durable Flag ... + /// + public bool QueueDurable { get; set; } = true; + + /// + /// Queue Exclusive Flag ... + /// + public bool QueueExclusive { get; set; } = false; + + /// + /// Queue AutoDelete Flag ... + /// + public bool QueueAutoDelete { get; set; } = false; + + /// + /// Enable Automatic Recovery on Connection Loss ... + /// + public bool AutomaticRecoveryEnabled { get; set; } = true; + + /// + /// Network Recovery Interval in Seconds ... + /// + public int NetworkRecoveryIntervalSeconds { get; set; } = 5; + + /// + /// Prefetch Count for Consumer ... + /// + public ushort PrefetchCount { get; set; } = 1; + + /// + /// Enable Manual Acknowledgement ... + /// + public bool ManualAck { get; set; } = true; + } +} \ No newline at end of file diff --git a/Providers/XKafkaConsumerServiceBase.cs b/Providers/XKafkaConsumerServiceBase.cs index 307dddb..5ed8520 100644 --- a/Providers/XKafkaConsumerServiceBase.cs +++ b/Providers/XKafkaConsumerServiceBase.cs @@ -1,9 +1,14 @@ using System; +using System.Threading; +using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; +using xCommons.Extensions; using xEventService.Configuration; using xEventService.Constants; +using xEventService.Extensions; +using xExceptions.Constants; namespace xEventService.Providers { @@ -15,18 +20,35 @@ namespace xEventService.Providers /// /// Consumed Topic Messages ... /// - private readonly string topic; + protected readonly string topic; + + /// + /// Check is Paused ... + /// + private bool paused = false; + + /// + /// Allowed Senders ... + /// + public readonly string[] allowedSenders; + + /// + /// Allowed Actions ... + /// + public readonly string[] allowedActions; /// /// Kafka Consumer Class ... /// - public readonly IConsumer consumer; + protected readonly IConsumer consumer; - private readonly XEventServiceConfiguration configuration; + protected readonly XEventServiceConfiguration configuration; - private readonly ILogger logger; + protected readonly ILogger logger; protected XKafkaConsumerServiceBase( + string[] allowedActions, + string[] allowedSenders, XEventServiceConfiguration configuration, ILogger logger, string topic = XEventServiceConstants.XDefaultTopic @@ -37,6 +59,24 @@ namespace xEventService.Providers this.logger = logger; this.configuration = configuration; + // + // Validate Actions ... + if (!allowedActions + .IsAllowedCollection()) + { + XException.InvalidArgs.Throw(); + } + this.allowedActions = [.. allowedActions]; + + // + // Validate Senders ... + if (!allowedSenders + .IsAllowedCollection()) + { + XException.InvalidArgs.Throw(); + } + this.allowedSenders = [.. allowedSenders]; + // // Prepare Kafka Consumer Configuration ... var config = new ConsumerConfig @@ -49,6 +89,148 @@ namespace xEventService.Providers consumer = new ConsumerBuilder(config).Build(); } + + /// + /// Pause Consuming ... + /// + /// + public bool Pause() + { + // + var result = !paused; + if (!result) + { + return result; + } + + // + paused = true; + result = paused; + return result; + } + + /// + /// Resume Consuming ... + /// + /// + public bool Resume() + { + // + var result = paused; + if (!result) + { + return result; + } + + // + paused = false; + result = !paused; + return result; + } + + /// + /// Default Action Execution for Background Services ... + /// + /// + /// + 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 ... + await consumer.Consume(async (consumeResult, cancellationToken) => + { + // Check Pausing ... + if (!paused) + { + // Logging Message Recieved ... + if (enableLogging) + { + // Encrypt Message ... + var expirationOffset = DateTimeOffset.UtcNow + .AddMinutes(encryptionLifetimeInMinute); + var encryptedMessage = cryptographyService.Encrypt( + text: consumeResult?.Message.Value ?? string.Empty, + expirationOffset: expirationOffset + ); + logger + .LogInformation( + "Consume Recieved Message from Kafka Topic: {Topic}, Offset: {Offset}, Partition: {Partition}, Message: {Message}", + topic, + consumeResult?.Offset, + consumeResult?.Partition, + encryptedMessage + ); + } + + // Produce Message Model ... + var model = consumeResult + .FillConsumeResult() + ?? throw RpkKafkaExceptionHelper + .GetInvalidKafkaConsumeResultException(); + + // Check Sender and Actions Owning ... + var isOwned = model.IsOwned( + allowedActions, + allowedSenders + ); + if (isOwned) + { + // Notify Extended Classes for Consume Message ... + await ConsumeAsync( + model, + cancellationToken + ); + } + + // Commit Offset after Consuming ... + consumer?.Commit(consumeResult); + } + }, + cts.Token); + } + catch (ConsumeException ex) + { + // Log Exception ... + if (enableLogging) + { + logger.LogError( + ex, + "Kafka Message Consume Failed: {Error}", + ex.Error.Reason + ); + } + await Task.Delay(1000, cancellationToken); + } + catch (Exception ex) + { + // Log Exception ... + if (enableLogging) + { + logger.LogError( + ex, + "Exception: {Exc}", + ex.Message + ); + } + await Task.Delay(1000, cancellationToken); + } + } + } + /// /// Dispose Implementation ... /// diff --git a/Providers/XRabbitMQConsumerServiceBase.cs b/Providers/XRabbitMQConsumerServiceBase.cs new file mode 100644 index 0000000..1fa5b72 --- /dev/null +++ b/Providers/XRabbitMQConsumerServiceBase.cs @@ -0,0 +1,50 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using RabbitMQ.Client; +using xCommons.Extensions; +using xEventService.Configuration; +using xEventService.Extensions; +using xExceptions.Constants; + +namespace xEventService.Providers +{ + /// + /// an Abstract Kafka Consumer Service ... + /// + public class XRabbitMQConsumerServiceBase : BackgroundService + { + protected readonly IModel channel; + protected readonly IConnection connection; + protected readonly XEventServiceConfiguration configuration; + protected readonly ILogger logger; + + public XRabbitMQConsumerServiceBase( + XEventServiceConfiguration configuration, + ILogger logger + ) + { + // + this.logger = logger; + this.configuration = configuration; + + // + // Validate ... + if (!configuration.IsValid() || + !configuration.RabbitMQ.IsValid() || + configuration.Broker != Constants.XEventBroker.XRabbitMQ + ) + { + XException.InvalidConfiguration.Throw(); + } + } + + public abstract Task HandleMessageAsync( + XEventMessage message, + CancellationToken cancellationToken = default + ); + } +} \ No newline at end of file diff --git a/Providers/XRabbitMQProducerService.cs b/Providers/XRabbitMQProducerService.cs new file mode 100644 index 0000000..74a7b80 --- /dev/null +++ b/Providers/XRabbitMQProducerService.cs @@ -0,0 +1,298 @@ +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 +{ + /// + /// RabbitMQ Producer Service Implementation ... + /// + public class XRabbitMQProducerService : IXRabbitMQProducerService + { + private IModel channel; + private IConnection connection; + private readonly object connectionLock = new(); + private readonly ILogger logger; + private readonly XEventServiceConfiguration configuration; + + public XRabbitMQProducerService( + ILogger logger, + XEventServiceConfiguration configuration + ) + { + // + this.logger = logger; + this.configuration = configuration; + + // + // Validate Configurations ... + if (!configuration.IsValid() || + !configuration.RabbitMQ.IsValid() || + configuration.Broker != XEventBroker.XRabbitMQ + ) + { + XException.InvalidConfiguration.Throw(); + } + + // + // Initialize Connection ... + CreateConnection(); + } + + /// + /// Sending Message with Full Request Descriptor ... + /// + public async Task 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; + } + + /// + /// Sending a Message to Exchange ... + /// + public async Task 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; + } + + /// + /// Dispose Object ... + /// + public void Dispose() + { + // + try + { + // + channel.Close(); + channel.Dispose(); + channel = null; + + // + connection.Close(); + connection.Dispose(); + connection = null; + } + catch + { } + + // + GC.SuppressFinalize(this); + } + + // + #region Private ... + /// + /// Create Connection to RabbitMQ Server ... + /// + 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(); + } + + /// + /// Check Connection is Alive ... + /// + private bool IsConnected() + { + // + var result = + !connection.IsNullOrDefault() && + connection.IsOpen && + !channel.IsNull() && + channel.IsOpen; + + // + return result; + } + + /// + /// Ensure Connection Established ... + /// + public async Task EnsureConnectionAsync() + { + // + try + { + await EnsureChannelAsync(); + return true; + } + catch + { + return false; + } + } + + /// + /// Ensure Channel is Open and Exchange/Queue Declared ... + /// + 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() + ); + } + + // + await Task.CompletedTask; + } + #endregion + } +} \ No newline at end of file