try to refactor Kafka and add supports for RabbitMQ ...

This commit is contained in:
2026-08-18 17:53:35 +03:30
parent 3e0bb2c25f
commit b20c2f8db7
9 changed files with 861 additions and 5 deletions
+186 -4
View File
@@ -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
/// <summary>
/// Consumed Topic Messages ...
/// </summary>
private readonly string topic;
protected readonly string topic;
/// <summary>
/// Check is Paused ...
/// </summary>
private bool paused = false;
/// <summary>
/// Allowed Senders ...
/// </summary>
public readonly string[] allowedSenders;
/// <summary>
/// Allowed Actions ...
/// </summary>
public readonly string[] allowedActions;
/// <summary>
/// Kafka Consumer Class ...
/// </summary>
public readonly IConsumer<Ignore, string> consumer;
protected readonly IConsumer<Ignore, string> consumer;
private readonly XEventServiceConfiguration configuration;
protected readonly XEventServiceConfiguration configuration;
private readonly ILogger<XKafkaConsumerServiceBase> logger;
protected readonly ILogger<XKafkaConsumerServiceBase> logger;
protected XKafkaConsumerServiceBase(
string[] allowedActions,
string[] allowedSenders,
XEventServiceConfiguration configuration,
ILogger<XKafkaConsumerServiceBase> 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<Ignore, string>(config).Build();
}
/// <summary>
/// Pause Consuming ...
/// </summary>
/// <returns></returns>
public bool Pause()
{
//
var result = !paused;
if (!result)
{
return result;
}
//
paused = true;
result = paused;
return result;
}
/// <summary>
/// Resume Consuming ...
/// </summary>
/// <returns></returns>
public bool Resume()
{
//
var result = paused;
if (!result)
{
return result;
}
//
paused = false;
result = !paused;
return result;
}
/// <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 ...
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);
}
}
}
/// <summary>
/// Dispose Implementation ...
/// </summary>
+50
View File
@@ -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
{
/// <summary>
/// an Abstract Kafka Consumer Service ...
/// </summary>
public class XRabbitMQConsumerServiceBase : BackgroundService
{
protected readonly IModel channel;
protected readonly IConnection connection;
protected readonly XEventServiceConfiguration configuration;
protected readonly ILogger<XRabbitMQConsumerServiceBase> logger;
public XRabbitMQConsumerServiceBase(
XEventServiceConfiguration configuration,
ILogger<XRabbitMQConsumerServiceBase> 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
);
}
}
+298
View File
@@ -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
{
/// <summary>
/// RabbitMQ Producer Service Implementation ...
/// </summary>
public class XRabbitMQProducerService : IXRabbitMQProducerService
{
private IModel channel;
private IConnection connection;
private readonly object connectionLock = new();
private readonly ILogger<XRabbitMQProducerService> logger;
private readonly XEventServiceConfiguration configuration;
public XRabbitMQProducerService(
ILogger<XRabbitMQProducerService> 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();
}
/// <summary>
/// Sending Message with Full Request Descriptor ...
/// </summary>
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 Exchange ...
/// </summary>
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>
public 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()
);
}
//
await Task.CompletedTask;
}
#endregion
}
}