From fe00b96a3a8415dc0791127692b13397ebc22661 Mon Sep 17 00:00:00 2001 From: Fabio Della Rosa Date: Wed, 17 Jun 2026 17:12:53 +0200 Subject: [PATCH 1/4] [ALPHA] - Kafka - Add Schema Registry integration, broker auth config, and Json/Avro/Protobuf wire serialization. --- .../Program.cs | 9 +- .../Consumers/CommandConsumerBase.cs | 14 +- .../Consumers/ConsumerBase.cs | 22 +- .../Consumers/DomainEventsConsumerBase.cs | 26 +-- ...se.cs => IntegrationEventsConsumerBase.cs} | 25 +- .../Consumers/KafkaSubscriber.cs | 60 +++-- .../KafkaBrokerStarter.cs | 6 +- .../Models/KafkaConfiguration.cs | 79 ++++++- .../Muflone.Transport.Kafka.csproj | 30 ++- src/Muflone.Transport.Kafka/README.md | 38 +++ .../Serialization/IKafkaMessageSerializer.cs | 14 ++ .../Serialization/KafkaEnvelope.cs | 7 + .../LegacyKafkaMessageSerializer.cs | 35 +++ .../Protobuf/kafka_envelope.proto | 8 + .../SchemaRegistryKafkaMessageSerializer.cs | 220 ++++++++++++++++++ src/Muflone.Transport.Kafka/ServiceBus.cs | 20 +- 16 files changed, 519 insertions(+), 94 deletions(-) rename src/Muflone.Transport.Kafka/Consumers/{IntegrazionEventsConsumerBase.cs => IntegrationEventsConsumerBase.cs} (55%) create mode 100644 src/Muflone.Transport.Kafka/Serialization/IKafkaMessageSerializer.cs create mode 100644 src/Muflone.Transport.Kafka/Serialization/KafkaEnvelope.cs create mode 100644 src/Muflone.Transport.Kafka/Serialization/LegacyKafkaMessageSerializer.cs create mode 100644 src/Muflone.Transport.Kafka/Serialization/Protobuf/kafka_envelope.proto create mode 100644 src/Muflone.Transport.Kafka/Serialization/SchemaRegistryKafkaMessageSerializer.cs diff --git a/src/Muflone.Transport.Kafka.AppTests/Program.cs b/src/Muflone.Transport.Kafka.AppTests/Program.cs index 4c05c8d..31cee26 100644 --- a/src/Muflone.Transport.Kafka.AppTests/Program.cs +++ b/src/Muflone.Transport.Kafka.AppTests/Program.cs @@ -11,7 +11,12 @@ HostApplicationBuilder builder = Host.CreateApplicationBuilder(args); var kafkaConfiguration = new KafkaConfiguration("localhost:9092", "test-group", "test-group"); - +kafkaConfiguration.BrokerUsername = "kafka-user"; +kafkaConfiguration.BrokerPassword = "kafka-password"; +kafkaConfiguration.SchemaRegistryUrl = "https://localhost:8081"; +kafkaConfiguration.SchemaRegistryUsername = "schema-user"; +kafkaConfiguration.SchemaRegistryPassword = "schema-password"; +kafkaConfiguration.SerializationFormat = KafkaSerializationFormat.Json; builder.Services.AddMufloneTransportKafka(new NullLoggerFactory(), kafkaConfiguration); builder.Services.AddHostedService(); @@ -20,4 +25,4 @@ using IHost host = builder.Build(); -await host.RunAsync(); \ No newline at end of file +await host.RunAsync(); diff --git a/src/Muflone.Transport.Kafka/Consumers/CommandConsumerBase.cs b/src/Muflone.Transport.Kafka/Consumers/CommandConsumerBase.cs index 5dc8b12..fe5aac7 100644 --- a/src/Muflone.Transport.Kafka/Consumers/CommandConsumerBase.cs +++ b/src/Muflone.Transport.Kafka/Consumers/CommandConsumerBase.cs @@ -21,13 +21,11 @@ public Task ConsumeAsync(T message, CancellationToken cancellationToken = defaul return HandlerAsync.HandleAsync(message, cancellationToken); } - public Task StartAsync(CancellationToken cancellationToken = default) - { - return StartConsumerAsync(ConsumeAsync, cancellationToken); - } + public Task StartAsync(CancellationToken cancellationToken = default) => + StartConsumerAsync(ConsumeAsync, cancellationToken); + - public Task StopAsync(CancellationToken cancellationToken = default) - { - return StopConsumerAsync(cancellationToken); - } + public Task StopAsync(CancellationToken cancellationToken = default) => + StopConsumerAsync(cancellationToken); + } diff --git a/src/Muflone.Transport.Kafka/Consumers/ConsumerBase.cs b/src/Muflone.Transport.Kafka/Consumers/ConsumerBase.cs index 46a7090..b703ce3 100644 --- a/src/Muflone.Transport.Kafka/Consumers/ConsumerBase.cs +++ b/src/Muflone.Transport.Kafka/Consumers/ConsumerBase.cs @@ -3,6 +3,7 @@ using Muflone.Messages; using Muflone.Persistence; using Muflone.Transport.Kafka.Models; +using Muflone.Transport.Kafka.Serialization; namespace Muflone.Transport.Kafka.Consumers; @@ -10,9 +11,9 @@ public abstract class ConsumerBase : IAsyncDisposable { protected readonly KafkaConfiguration Configuration; protected readonly ILogger Logger; - protected readonly ISerializer MessageSerializer; + private readonly IKafkaMessageSerializer _kafkaMessageSerializer; - private IConsumer? _consumer; + private IConsumer? _consumer; private Task? _consumerTask; private CancellationTokenSource? _stoppingTokenSource; @@ -23,7 +24,7 @@ protected ConsumerBase( { Configuration = configuration ?? throw new ArgumentNullException(nameof(configuration)); Logger = loggerFactory?.CreateLogger(GetType()) ?? throw new ArgumentNullException(nameof(loggerFactory)); - MessageSerializer = messageSerializer ?? new Serializer(); + _kafkaMessageSerializer = KafkaMessageSerializerFactory.Create(configuration, messageSerializer ?? new Serializer()); } protected Task StartConsumerAsync( @@ -75,7 +76,7 @@ protected async Task StopConsumerAsync(CancellationToken cancellationToken = def } } - private IConsumer CreateConsumer() + private IConsumer CreateConsumer() where TMessage : class, IMessage { var consumerConfig = new ConsumerConfig @@ -88,11 +89,13 @@ private IConsumer CreateConsumer() AllowAutoCreateTopics = true }; - return new ConsumerBuilder(consumerConfig).Build(); + Configuration.ApplyBrokerAuthentication(consumerConfig); + + return new ConsumerBuilder(consumerConfig).Build(); } private async Task ConsumeLoopAsync( - IConsumer consumer, + IConsumer consumer, Func messageHandler, CancellationToken cancellationToken) where TMessage : class, IMessage @@ -101,7 +104,7 @@ private async Task ConsumeLoopAsync( { while (!cancellationToken.IsCancellationRequested) { - ConsumeResult? consumeResult; + ConsumeResult? consumeResult; try { @@ -118,8 +121,8 @@ private async Task ConsumeLoopAsync( try { - var message = await MessageSerializer - .DeserializeAsync(consumeResult.Message.Value, cancellationToken) + var message = await _kafkaMessageSerializer + .DeserializeAsync(consumeResult.Topic, consumeResult.Message.Value, cancellationToken) .ConfigureAwait(false); if (message == null) @@ -168,6 +171,7 @@ private async Task ConsumeLoopAsync( public async ValueTask DisposeAsync() { await StopConsumerAsync().ConfigureAwait(false); + await _kafkaMessageSerializer.DisposeAsync().ConfigureAwait(false); GC.SuppressFinalize(this); } } diff --git a/src/Muflone.Transport.Kafka/Consumers/DomainEventsConsumerBase.cs b/src/Muflone.Transport.Kafka/Consumers/DomainEventsConsumerBase.cs index 56b3ed2..8946567 100644 --- a/src/Muflone.Transport.Kafka/Consumers/DomainEventsConsumerBase.cs +++ b/src/Muflone.Transport.Kafka/Consumers/DomainEventsConsumerBase.cs @@ -5,19 +5,14 @@ namespace Muflone.Transport.Kafka.Consumers; -public abstract class DomainEventsConsumerBase : ConsumerBase, IDomainEventConsumer +public abstract class DomainEventsConsumerBase( + KafkaConfiguration configuration, + ILoggerFactory loggerFactory, + ISerializer? messageSerializer = null) : ConsumerBase(configuration, loggerFactory, messageSerializer), IDomainEventConsumer where T : DomainEvent { protected abstract IEnumerable> HandlersAsync { get; } - protected DomainEventsConsumerBase( - KafkaConfiguration configuration, - ILoggerFactory loggerFactory, - ISerializer? messageSerializer = null) - : base(configuration, loggerFactory, messageSerializer) - { - } - public async Task ConsumeAsync(T message, CancellationToken cancellationToken = default) { ArgumentNullException.ThrowIfNull(message); @@ -26,13 +21,10 @@ public async Task ConsumeAsync(T message, CancellationToken cancellationToken = await handlerAsync.HandleAsync(message, cancellationToken).ConfigureAwait(false); } - public Task StartAsync(CancellationToken cancellationToken = default) - { - return StartConsumerAsync(ConsumeAsync, cancellationToken); - } + public Task StartAsync(CancellationToken cancellationToken = default) => + StartConsumerAsync(ConsumeAsync, cancellationToken); - public Task StopAsync(CancellationToken cancellationToken = default) - { - return StopConsumerAsync(cancellationToken); - } + public Task StopAsync(CancellationToken cancellationToken = default) => + StopConsumerAsync(cancellationToken); + } diff --git a/src/Muflone.Transport.Kafka/Consumers/IntegrazionEventsConsumerBase.cs b/src/Muflone.Transport.Kafka/Consumers/IntegrationEventsConsumerBase.cs similarity index 55% rename from src/Muflone.Transport.Kafka/Consumers/IntegrazionEventsConsumerBase.cs rename to src/Muflone.Transport.Kafka/Consumers/IntegrationEventsConsumerBase.cs index 6bec072..77acd6c 100644 --- a/src/Muflone.Transport.Kafka/Consumers/IntegrazionEventsConsumerBase.cs +++ b/src/Muflone.Transport.Kafka/Consumers/IntegrationEventsConsumerBase.cs @@ -5,19 +5,15 @@ namespace Muflone.Transport.Kafka.Consumers; -public abstract class IntegrationEventsConsumerBase : ConsumerBase, IIntegrationEventConsumer +public abstract class IntegrationEventsConsumerBase( + KafkaConfiguration configuration, + ILoggerFactory loggerFactory, + ISerializer? messageSerializer = null) + : ConsumerBase(configuration, loggerFactory, messageSerializer), IIntegrationEventConsumer where T : IntegrationEvent { protected abstract IEnumerable> HandlersAsync { get; } - protected IntegrationEventsConsumerBase( - KafkaConfiguration configuration, - ILoggerFactory loggerFactory, - ISerializer? messageSerializer = null) - : base(configuration, loggerFactory, messageSerializer) - { - } - public async Task ConsumeAsync(T message, CancellationToken cancellationToken = default) { ArgumentNullException.ThrowIfNull(message); @@ -37,14 +33,3 @@ public Task StopAsync(CancellationToken cancellationToken = default) } } -public abstract class IntegrazionEventsConsumerBase : IntegrationEventsConsumerBase - where T : IntegrationEvent -{ - protected IntegrazionEventsConsumerBase( - KafkaConfiguration configuration, - ILoggerFactory loggerFactory, - ISerializer? messageSerializer = null) - : base(configuration, loggerFactory, messageSerializer) - { - } -} diff --git a/src/Muflone.Transport.Kafka/Consumers/KafkaSubscriber.cs b/src/Muflone.Transport.Kafka/Consumers/KafkaSubscriber.cs index 6bb88a9..1af1a0b 100644 --- a/src/Muflone.Transport.Kafka/Consumers/KafkaSubscriber.cs +++ b/src/Muflone.Transport.Kafka/Consumers/KafkaSubscriber.cs @@ -1,11 +1,13 @@ using Confluent.Kafka; using Microsoft.Extensions.Logging; using Muflone.Messages; +using Muflone.Persistence; using Muflone.Transport.Kafka.Models; +using Muflone.Transport.Kafka.Serialization; +using KafkaConsumerClient = Confluent.Kafka.IConsumer; namespace Muflone.Transport.Kafka.Consumers; -using KafkaConsumerClient = Confluent.Kafka.IConsumer; public class KafkaSubscriber( ILoggerFactory loggerFactory, @@ -13,13 +15,14 @@ public class KafkaSubscriber( KafkaConfiguration configuration) : MessageSubscriberBase(loggerFactory, serviceProvider) { private readonly ILogger _logger = loggerFactory.CreateLogger(); + private readonly IKafkaMessageSerializer _kafkaMessageSerializer = KafkaMessageSerializerFactory.Create(configuration, new Serializer()); protected override async Task StopChannelAsync(HandlerSubscription handlerSubscription) { - if (handlerSubscription.Channel == null) + if (handlerSubscription.Channel is null) return; - await handlerSubscription.Channel.StopAsync().ConfigureAwait(false); + await handlerSubscription.Channel.StopAsync(); handlerSubscription.Channel = null; } @@ -38,7 +41,9 @@ protected override Task InitChannelAsync(HandlerSubscription(consumerConfig).Build(); + configuration.ApplyBrokerAuthentication(consumerConfig); + + var consumer = new ConsumerBuilder(consumerConfig).Build(); consumer.Subscribe(topicName); _logger.LogInformation( @@ -55,14 +60,14 @@ protected override Task InitSubscriptionAsync(HandlerSubscription { while (!channel.StoppingTokenSource.IsCancellationRequested) { - ConsumeResult? consumeResult; + ConsumeResult? consumeResult; try { @@ -78,12 +83,19 @@ protected override Task InitSubscriptionAsync(HandlerSubscription handlerSubscription) - { - return handlerSubscription.EventTypeName.ToLowerInvariant(); - } + private static string GetTopicName(HandlerSubscription handlerSubscription) + => handlerSubscription.EventTypeName.ToLowerInvariant(); + + // topicHandler-{Guid} + // private string GetGroupId(HandlerSubscription handlerSubscription) + // { + // var groupId = $"{configuration.GroupId}.{handlerSubscription.EventTypeName}"; + // if (groupId.EndsWith("Consumer", StringComparison.InvariantCultureIgnoreCase)) + // groupId = groupId[..^"Consumer".Length]; + // + // if (!handlerSubscription.IsCommandHandler || !handlerSubscription.IsSingletonHandler) + // groupId = $"{groupId}.{handlerSubscription.HandlerSubscriptionId}"; + // + // return groupId; + // } + + // topicHandler private string GetGroupId(HandlerSubscription handlerSubscription) { var groupId = $"{configuration.GroupId}.{handlerSubscription.EventTypeName}"; + if (groupId.EndsWith("Consumer", StringComparison.InvariantCultureIgnoreCase)) groupId = groupId[..^"Consumer".Length]; - if (!handlerSubscription.IsCommandHandler || !handlerSubscription.IsSingletonHandler) - groupId = $"{groupId}.{handlerSubscription.HandlerSubscriptionId}"; + if (handlerSubscription.Configuration?.InstanceId is { Length: > 0 } instanceId) + groupId = $"{groupId}.{instanceId}"; return groupId; } @@ -126,13 +152,13 @@ public sealed class KafkaSubscriptionChannel(KafkaConsumerClient consumer) public async Task StopAsync() { - await StoppingTokenSource.CancelAsync().ConfigureAwait(false); + await StoppingTokenSource.CancelAsync(); - if (ConsumerTask != null) + if (ConsumerTask is not null) { try { - await ConsumerTask.ConfigureAwait(false); + await ConsumerTask; } catch (OperationCanceledException) { diff --git a/src/Muflone.Transport.Kafka/KafkaBrokerStarter.cs b/src/Muflone.Transport.Kafka/KafkaBrokerStarter.cs index 3197ea3..3302901 100644 --- a/src/Muflone.Transport.Kafka/KafkaBrokerStarter.cs +++ b/src/Muflone.Transport.Kafka/KafkaBrokerStarter.cs @@ -10,8 +10,6 @@ public async Task StartAsync(CancellationToken cancellationToken) await consumer.StartAsync(cancellationToken).ConfigureAwait(false); } - public Task StopAsync(CancellationToken cancellationToken) - { - return Task.CompletedTask; - } + public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; + } diff --git a/src/Muflone.Transport.Kafka/Models/KafkaConfiguration.cs b/src/Muflone.Transport.Kafka/Models/KafkaConfiguration.cs index 572ea18..47d23d1 100644 --- a/src/Muflone.Transport.Kafka/Models/KafkaConfiguration.cs +++ b/src/Muflone.Transport.Kafka/Models/KafkaConfiguration.cs @@ -1,14 +1,34 @@ using Confluent.Kafka; +using Confluent.SchemaRegistry; using System.Globalization; namespace Muflone.Transport.Kafka.Models; +public enum KafkaSerializationFormat +{ + Legacy = 0, + Json = 1, + Avro = 2, + Protobuf = 3 +} + public class KafkaConfiguration { public string BootstrapServers { get; } public string GroupId { get; } public string ClientId { get; } public AutoOffsetReset AutoOffsetReset { get; } + public KafkaSerializationFormat SerializationFormat { get; set; } = KafkaSerializationFormat.Legacy; + public SecurityProtocol? BrokerSecurityProtocol { get; set; } + public SaslMechanism? BrokerSaslMechanism { get; set; } + public string? BrokerUsername { get; set; } + public string? BrokerPassword { get; set; } + public string? SchemaRegistryUrl { get; set; } + public string? SchemaRegistryUsername { get; set; } + public string? SchemaRegistryPassword { get; set; } + public AuthCredentialsSource SchemaRegistryBasicAuthCredentialsSource { get; set; } = AuthCredentialsSource.UserInfo; + public int? SchemaRegistryMaxCachedSchemas { get; set; } + public bool SchemaRegistryAutoRegisterSchemas { get; set; } = true; public Func TopicNamingConvention { get; set; } = type => type.Name.ToLower(CultureInfo.InvariantCulture); @@ -26,7 +46,18 @@ public KafkaConfiguration( string bootstrapServers, string groupId, string clientId, - AutoOffsetReset autoOffsetReset = AutoOffsetReset.Earliest) + AutoOffsetReset autoOffsetReset = AutoOffsetReset.Earliest, + KafkaSerializationFormat serializationFormat = default, + SecurityProtocol? brokerSecurityProtocol = null, + SaslMechanism? brokerSaslMechanism = null, + string? brokerUsername = null, + string? brokerPassword = null, + string? schemaRegistryUrl = null, + string? schemaRegistryUsername = null, + string? schemaRegistryPassword = null, + AuthCredentialsSource schemaRegistryBasicAuthCredentialsSource = default, + int? schemaRegistryMaxCachedSchemas = null, + bool schemaRegistryAutoRegisterSchemas = false) { BootstrapServers = string.IsNullOrWhiteSpace(bootstrapServers) ? throw new ArgumentNullException(nameof(bootstrapServers)) @@ -36,6 +67,17 @@ public KafkaConfiguration( : groupId; ClientId = string.IsNullOrWhiteSpace(clientId) ? groupId : clientId; AutoOffsetReset = autoOffsetReset; + SerializationFormat = serializationFormat; + BrokerSecurityProtocol = brokerSecurityProtocol; + BrokerSaslMechanism = brokerSaslMechanism; + BrokerUsername = brokerUsername; + BrokerPassword = brokerPassword; + SchemaRegistryUrl = schemaRegistryUrl; + SchemaRegistryUsername = schemaRegistryUsername; + SchemaRegistryPassword = schemaRegistryPassword; + SchemaRegistryBasicAuthCredentialsSource = schemaRegistryBasicAuthCredentialsSource; + SchemaRegistryMaxCachedSchemas = schemaRegistryMaxCachedSchemas; + SchemaRegistryAutoRegisterSchemas = schemaRegistryAutoRegisterSchemas; } public string GetTopicName(Type messageType) @@ -43,4 +85,39 @@ public string GetTopicName(Type messageType) ArgumentNullException.ThrowIfNull(messageType); return TopicNamingConvention(messageType); } + + internal void ApplyBrokerAuthentication(ClientConfig clientConfig) + { + ArgumentNullException.ThrowIfNull(clientConfig); + + if (string.IsNullOrWhiteSpace(BrokerUsername) || string.IsNullOrWhiteSpace(BrokerPassword)) + return; + + clientConfig.SecurityProtocol = BrokerSecurityProtocol ?? SecurityProtocol.SaslSsl; + clientConfig.SaslMechanism = BrokerSaslMechanism ?? SaslMechanism.Plain; + clientConfig.SaslUsername = BrokerUsername; + clientConfig.SaslPassword = BrokerPassword; + } + + internal SchemaRegistryConfig ToSchemaRegistryConfig() + { + if (string.IsNullOrWhiteSpace(SchemaRegistryUrl)) + throw new InvalidOperationException("SchemaRegistryUrl is required when using Json, Avro or Protobuf serialization."); + + var config = new SchemaRegistryConfig + { + Url = SchemaRegistryUrl + }; + + if (SchemaRegistryMaxCachedSchemas.HasValue) + config.MaxCachedSchemas = SchemaRegistryMaxCachedSchemas.Value; + + if (!string.IsNullOrWhiteSpace(SchemaRegistryUsername) && !string.IsNullOrWhiteSpace(SchemaRegistryPassword)) + { + config.BasicAuthCredentialsSource = SchemaRegistryBasicAuthCredentialsSource; + config.BasicAuthUserInfo = $"{SchemaRegistryUsername}:{SchemaRegistryPassword}"; + } + + return config; + } } diff --git a/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj b/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj index 00d9dd3..2d0fd4e 100644 --- a/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj +++ b/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj @@ -6,7 +6,7 @@ enable Fabio Della Rosa https://github.com/CQRS-Muflone/Muflone.Transport.Kafka - Logo.png + README.md Git Broker, Kafka, CQRS, Saga, Event Sourcing @@ -21,10 +21,10 @@ - + True \ @@ -32,11 +32,25 @@ - - - - - + + + + + + + + all + runtime; build; native; contentfiles; analyzers; buildtransitive + + + + + + + + + + diff --git a/src/Muflone.Transport.Kafka/README.md b/src/Muflone.Transport.Kafka/README.md index f7ce9d9..63df329 100644 --- a/src/Muflone.Transport.Kafka/README.md +++ b/src/Muflone.Transport.Kafka/README.md @@ -22,6 +22,8 @@ Explicit consumers are still supported, but the transport can now also dispatch - **Automatic Handler Dispatch** - Supports `IMessageSubscriber` and `MessageHandlersStarter`, aligned with the RabbitMQ transport - **Explicit Consumers** - Supports `CommandConsumerBase`, `DomainEventsConsumerBase`, and `IntegrationEventsConsumerBase` - **Topic Naming Convention** - Customize topic naming through `KafkaConfiguration.TopicNamingConvention` +- **Schema Registry Integration** - Supports Confluent Schema Registry for `Json`, `Avro`, and `Protobuf` +- **Basic Auth** - Supports SASL basic auth for Kafka brokers and HTTP Basic Auth for Schema Registry - **Async-first** - All operations use async/await - **.NET 10** - Built for `net10.0` with nullable reference types @@ -191,6 +193,13 @@ var kafkaConfiguration = new KafkaConfiguration( autoOffsetReset: AutoOffsetReset.Earliest ); +kafkaConfiguration.BrokerUsername = "kafka-user"; +kafkaConfiguration.BrokerPassword = "kafka-password"; +kafkaConfiguration.SchemaRegistryUrl = "https://localhost:8081"; +kafkaConfiguration.SchemaRegistryUsername = "schema-user"; +kafkaConfiguration.SchemaRegistryPassword = "schema-password"; +kafkaConfiguration.SerializationFormat = KafkaSerializationFormat.Json; + // Register handlers builder.Services.AddScoped, CreateOrderCommandHandler>(); builder.Services.AddScoped, OrderCreatedEventHandler>(); @@ -275,6 +284,35 @@ new KafkaConfiguration(bootstrapServers, groupId, clientId, AutoOffsetReset.Earl | `clientId` | Client id used by the producer and subscribers | `groupId` | | `autoOffsetReset` | Kafka offset reset strategy | `Earliest` | +### Broker basic auth + +```csharp +kafkaConfiguration.BrokerUsername = "kafka-user"; +kafkaConfiguration.BrokerPassword = "kafka-password"; +kafkaConfiguration.BrokerSecurityProtocol = SecurityProtocol.SaslSsl; +kafkaConfiguration.BrokerSaslMechanism = SaslMechanism.Plain; +``` + +### Schema Registry and serialization format + +The transport keeps using the configured Muflone `ISerializer` for the domain message payload and wraps that payload in a Kafka envelope managed through Schema Registry. + +This means: + +- existing Muflone messages do not need to become Avro or Protobuf classes +- Kafka producer/consumer wire format is controlled by `KafkaSerializationFormat` +- Schema Registry subjects are registered through the Confluent .NET serdes packages + +```csharp +kafkaConfiguration.SchemaRegistryUrl = "https://localhost:8081"; +kafkaConfiguration.SchemaRegistryUsername = "schema-user"; +kafkaConfiguration.SchemaRegistryPassword = "schema-password"; +kafkaConfiguration.SchemaRegistryAutoRegisterSchemas = true; +kafkaConfiguration.SerializationFormat = KafkaSerializationFormat.Avro; +// or KafkaSerializationFormat.Json +// or KafkaSerializationFormat.Protobuf +``` + ### Custom topic naming diff --git a/src/Muflone.Transport.Kafka/Serialization/IKafkaMessageSerializer.cs b/src/Muflone.Transport.Kafka/Serialization/IKafkaMessageSerializer.cs new file mode 100644 index 0000000..9708846 --- /dev/null +++ b/src/Muflone.Transport.Kafka/Serialization/IKafkaMessageSerializer.cs @@ -0,0 +1,14 @@ +using Muflone.Messages; + +namespace Muflone.Transport.Kafka.Serialization; + +internal interface IKafkaMessageSerializer : IAsyncDisposable +{ + Task SerializeAsync(string topicName, TMessage message, CancellationToken cancellationToken) + where TMessage : class, IMessage; + + Task DeserializeAsync(string topicName, ReadOnlyMemory payload, CancellationToken cancellationToken) + where TMessage : class, IMessage; + + Task DeserializePayloadAsync(string topicName, ReadOnlyMemory payload, CancellationToken cancellationToken); +} diff --git a/src/Muflone.Transport.Kafka/Serialization/KafkaEnvelope.cs b/src/Muflone.Transport.Kafka/Serialization/KafkaEnvelope.cs new file mode 100644 index 0000000..3fdb616 --- /dev/null +++ b/src/Muflone.Transport.Kafka/Serialization/KafkaEnvelope.cs @@ -0,0 +1,7 @@ +namespace Muflone.Transport.Kafka.Serialization; + +internal sealed class KafkaEnvelope +{ + public string MessageType { get; init; } = string.Empty; + public string Payload { get; init; } = string.Empty; +} diff --git a/src/Muflone.Transport.Kafka/Serialization/LegacyKafkaMessageSerializer.cs b/src/Muflone.Transport.Kafka/Serialization/LegacyKafkaMessageSerializer.cs new file mode 100644 index 0000000..bc71f29 --- /dev/null +++ b/src/Muflone.Transport.Kafka/Serialization/LegacyKafkaMessageSerializer.cs @@ -0,0 +1,35 @@ +using System.Text; +using Muflone.Messages; +using Muflone.Persistence; + +namespace Muflone.Transport.Kafka.Serialization; + +internal sealed class LegacyKafkaMessageSerializer(ISerializer messageSerializer) : IKafkaMessageSerializer +{ + public async Task SerializeAsync(string topicName, TMessage message, CancellationToken cancellationToken) + where TMessage : class, IMessage + { + ArgumentNullException.ThrowIfNull(message); + + var serializedMessage = await messageSerializer.SerializeAsync(message, cancellationToken); + return Encoding.UTF8.GetBytes(serializedMessage); + } + + public async Task DeserializeAsync(string topicName, ReadOnlyMemory payload, CancellationToken cancellationToken) + where TMessage : class, IMessage + { + var serializedMessage = Encoding.UTF8.GetString(payload.Span); + return await messageSerializer.DeserializeAsync(serializedMessage, cancellationToken); + } + + public Task DeserializePayloadAsync(string topicName, ReadOnlyMemory payload, CancellationToken cancellationToken) + { + var serializedMessage = Encoding.UTF8.GetString(payload.Span); + return Task.FromResult(serializedMessage); + } + + public ValueTask DisposeAsync() + { + return ValueTask.CompletedTask; + } +} diff --git a/src/Muflone.Transport.Kafka/Serialization/Protobuf/kafka_envelope.proto b/src/Muflone.Transport.Kafka/Serialization/Protobuf/kafka_envelope.proto new file mode 100644 index 0000000..0779f1e --- /dev/null +++ b/src/Muflone.Transport.Kafka/Serialization/Protobuf/kafka_envelope.proto @@ -0,0 +1,8 @@ +syntax = "proto3"; + +option csharp_namespace = "Muflone.Transport.Kafka.Serialization.Protobuf"; + +message KafkaEnvelopeContract { + string message_type = 1; + string payload = 2; +} diff --git a/src/Muflone.Transport.Kafka/Serialization/SchemaRegistryKafkaMessageSerializer.cs b/src/Muflone.Transport.Kafka/Serialization/SchemaRegistryKafkaMessageSerializer.cs new file mode 100644 index 0000000..8ac885d --- /dev/null +++ b/src/Muflone.Transport.Kafka/Serialization/SchemaRegistryKafkaMessageSerializer.cs @@ -0,0 +1,220 @@ +using Avro; +using Avro.Generic; +using Confluent.Kafka; +using Confluent.SchemaRegistry; +using Confluent.SchemaRegistry.Serdes; +using Muflone.Persistence; +using Muflone.Transport.Kafka.Models; +using Muflone.Transport.Kafka.Serialization.Protobuf; +using JsonDeserializer = Confluent.SchemaRegistry.Serdes.JsonDeserializer; +using JsonSerializer = Confluent.SchemaRegistry.Serdes.JsonSerializer; +using ProtobufDeserializer = Confluent.SchemaRegistry.Serdes.ProtobufDeserializer; +using ProtobufSerializer = Confluent.SchemaRegistry.Serdes.ProtobufSerializer; + +namespace Muflone.Transport.Kafka.Serialization; + +internal static class KafkaMessageSerializerFactory +{ + public static IKafkaMessageSerializer Create(KafkaConfiguration configuration, ISerializer messageSerializer) + { + ArgumentNullException.ThrowIfNull(configuration); + ArgumentNullException.ThrowIfNull(messageSerializer); + + return configuration.SerializationFormat switch + { + KafkaSerializationFormat.Legacy => new LegacyKafkaMessageSerializer(messageSerializer), + KafkaSerializationFormat.Json or KafkaSerializationFormat.Avro or KafkaSerializationFormat.Protobuf + => new SchemaRegistryKafkaMessageSerializer(configuration, messageSerializer), + _ => throw new NotSupportedException($"Serialization format '{configuration.SerializationFormat}' is not supported.") + }; + } +} + +internal sealed class SchemaRegistryKafkaMessageSerializer : IKafkaMessageSerializer +{ + private static readonly Avro.Schema AvroEnvelopeSchema = Avro.Schema.Parse(""" + { + "type": "record", + "name": "KafkaEnvelope", + "namespace": "Muflone.Transport.Kafka.Serialization", + "fields": [ + { "name": "messageType", "type": "string" }, + { "name": "payload", "type": "string" } + ] + } + """); + + private readonly KafkaConfiguration _configuration; + private readonly ISerializer _messageSerializer; + private readonly CachedSchemaRegistryClient _schemaRegistryClient; + private readonly JsonSerializer? _jsonSerializer; + private readonly JsonDeserializer? _jsonDeserializer; + private readonly AvroSerializer? _avroSerializer; + private readonly AvroDeserializer? _avroDeserializer; + private readonly ProtobufSerializer? _protobufSerializer; + private readonly ProtobufDeserializer? _protobufDeserializer; + + public SchemaRegistryKafkaMessageSerializer(KafkaConfiguration configuration, ISerializer messageSerializer) + { + _configuration = configuration ?? throw new ArgumentNullException(nameof(configuration)); + _messageSerializer = messageSerializer ?? throw new ArgumentNullException(nameof(messageSerializer)); + + _schemaRegistryClient = new CachedSchemaRegistryClient(configuration.ToSchemaRegistryConfig()); + + switch (_configuration.SerializationFormat) + { + case KafkaSerializationFormat.Json: + _jsonSerializer = new JsonSerializer( + _schemaRegistryClient, + new JsonSerializerConfig + { + AutoRegisterSchemas = _configuration.SchemaRegistryAutoRegisterSchemas + }); + _jsonDeserializer = new JsonDeserializer(_schemaRegistryClient); + break; + case KafkaSerializationFormat.Avro: + _avroSerializer = new AvroSerializer( + _schemaRegistryClient, + new AvroSerializerConfig + { + AutoRegisterSchemas = _configuration.SchemaRegistryAutoRegisterSchemas + }); + _avroDeserializer = new AvroDeserializer(_schemaRegistryClient); + break; + case KafkaSerializationFormat.Protobuf: + _protobufSerializer = new ProtobufSerializer( + _schemaRegistryClient, + new ProtobufSerializerConfig + { + AutoRegisterSchemas = _configuration.SchemaRegistryAutoRegisterSchemas + }); + _protobufDeserializer = new ProtobufDeserializer(_schemaRegistryClient); + break; + default: + throw new NotSupportedException($"Serialization format '{_configuration.SerializationFormat}' is not supported."); + } + } + + public async Task SerializeAsync(string topicName, TMessage message, CancellationToken cancellationToken) + where TMessage : class, Muflone.Messages.IMessage + { + ArgumentNullException.ThrowIfNull(message); + + var payload = await _messageSerializer.SerializeAsync(message, cancellationToken); + var envelope = new KafkaEnvelope + { + MessageType = message.GetType().AssemblyQualifiedName ?? message.GetType().FullName ?? message.GetType().Name, + Payload = payload + }; + + var context = CreateSerializationContext(topicName); + + return _configuration.SerializationFormat switch + { + KafkaSerializationFormat.Json => await _jsonSerializer!.SerializeAsync(envelope, context), + KafkaSerializationFormat.Avro => await _avroSerializer!.SerializeAsync(ToAvroEnvelope(envelope), context), + KafkaSerializationFormat.Protobuf => await _protobufSerializer!.SerializeAsync(ToProtobufEnvelope(envelope), context), + _ => throw new NotSupportedException($"Serialization format '{_configuration.SerializationFormat}' is not supported.") + }; + } + + public async Task DeserializeAsync(string topicName, ReadOnlyMemory payload, CancellationToken cancellationToken) + where TMessage : class, Muflone.Messages.IMessage + { + var envelope = await DeserializeEnvelopeAsync(topicName, payload).ConfigureAwait(false); + if (envelope == null) + return default; + + return await _messageSerializer.DeserializeAsync(envelope.Payload, cancellationToken).ConfigureAwait(false); + } + + public async Task DeserializePayloadAsync(string topicName, ReadOnlyMemory payload, CancellationToken cancellationToken) + { + var envelope = await DeserializeEnvelopeAsync(topicName, payload).ConfigureAwait(false); + return envelope?.Payload; + } + + public async ValueTask DisposeAsync() + { + switch (_jsonSerializer) + { + case IAsyncDisposable asyncDisposable: + await asyncDisposable.DisposeAsync().ConfigureAwait(false); + break; + case IDisposable disposable: + disposable.Dispose(); + break; + } + + switch (_jsonDeserializer) + { + case IAsyncDisposable asyncDisposable: + await asyncDisposable.DisposeAsync().ConfigureAwait(false); + break; + case IDisposable disposable: + disposable.Dispose(); + break; + } + + _schemaRegistryClient.Dispose(); + } + + private async Task DeserializeEnvelopeAsync(string topicName, ReadOnlyMemory payload) + { + var context = CreateSerializationContext(topicName); + + return _configuration.SerializationFormat switch + { + KafkaSerializationFormat.Json => await _jsonDeserializer!.DeserializeAsync(payload, false, context).ConfigureAwait(false), + KafkaSerializationFormat.Avro => FromAvroEnvelope(await _avroDeserializer!.DeserializeAsync(payload, false, context).ConfigureAwait(false)), + KafkaSerializationFormat.Protobuf => FromProtobufEnvelope(await _protobufDeserializer!.DeserializeAsync(payload, false, context).ConfigureAwait(false)), + _ => throw new NotSupportedException($"Serialization format '{_configuration.SerializationFormat}' is not supported.") + }; + } + + private static SerializationContext CreateSerializationContext(string topicName) + { + return new SerializationContext(MessageComponentType.Value, topicName); + } + + private static GenericRecord ToAvroEnvelope(KafkaEnvelope envelope) + { + var record = new GenericRecord((RecordSchema)AvroEnvelopeSchema); + record.Add("messageType", envelope.MessageType); + record.Add("payload", envelope.Payload); + return record; + } + + private static KafkaEnvelope? FromAvroEnvelope(GenericRecord? envelope) + { + if (envelope == null) + return null; + + return new KafkaEnvelope + { + MessageType = envelope["messageType"]?.ToString() ?? string.Empty, + Payload = envelope["payload"]?.ToString() ?? string.Empty + }; + } + + private static KafkaEnvelopeContract ToProtobufEnvelope(KafkaEnvelope envelope) + { + return new KafkaEnvelopeContract + { + MessageType = envelope.MessageType, + Payload = envelope.Payload + }; + } + + private static KafkaEnvelope? FromProtobufEnvelope(KafkaEnvelopeContract? envelope) + { + if (envelope == null) + return null; + + return new KafkaEnvelope + { + MessageType = envelope.MessageType, + Payload = envelope.Payload + }; + } +} diff --git a/src/Muflone.Transport.Kafka/ServiceBus.cs b/src/Muflone.Transport.Kafka/ServiceBus.cs index 1d7fe12..df64269 100644 --- a/src/Muflone.Transport.Kafka/ServiceBus.cs +++ b/src/Muflone.Transport.Kafka/ServiceBus.cs @@ -5,6 +5,7 @@ using Muflone.Messages.Events; using Muflone.Persistence; using Muflone.Transport.Kafka.Models; +using Muflone.Transport.Kafka.Serialization; using System.Text; namespace Muflone.Transport.Kafka; @@ -12,10 +13,10 @@ namespace Muflone.Transport.Kafka; public sealed class ServiceBus : IServiceBus, IEventBus, IDisposable { private readonly KafkaConfiguration _configuration; - private readonly ISerializer _messageSerializer; + private readonly IKafkaMessageSerializer _kafkaMessageSerializer; private readonly ILogger _logger; private readonly object _producerLock = new(); - private IProducer? _producer; + private IProducer? _producer; public ServiceBus( KafkaConfiguration configuration, @@ -23,7 +24,7 @@ public ServiceBus( ISerializer? messageSerializer = null) { _configuration = configuration ?? throw new ArgumentNullException(nameof(configuration)); - _messageSerializer = messageSerializer ?? new Serializer(); + _kafkaMessageSerializer = KafkaMessageSerializerFactory.Create(configuration, messageSerializer ?? new Serializer()); _logger = loggerFactory?.CreateLogger(GetType()) ?? throw new ArgumentNullException(nameof(loggerFactory)); } @@ -44,8 +45,8 @@ public Task SendAsync(T command, CancellationToken cancellationToken = defaul private async Task ProduceAsync(TMessage message, CancellationToken cancellationToken) where TMessage : class, IMessage { - var serializedMessage = await _messageSerializer.SerializeAsync(message, cancellationToken).ConfigureAwait(false); var topicName = _configuration.GetTopicName(message.GetType()); + var serializedMessage = await _kafkaMessageSerializer.SerializeAsync(topicName, message, cancellationToken).ConfigureAwait(false); _logger.LogInformation( "Publishing message '{MessageId}' of type '{MessageType}' to topic '{TopicName}'", @@ -53,7 +54,7 @@ private async Task ProduceAsync(TMessage message, CancellationToken ca message.GetType().FullName, topicName); - var kafkaMessage = new Message + var kafkaMessage = new Message { Key = message.MessageId.ToString(), Value = serializedMessage, @@ -63,7 +64,7 @@ private async Task ProduceAsync(TMessage message, CancellationToken ca await GetProducer().ProduceAsync(topicName, kafkaMessage, cancellationToken).ConfigureAwait(false); } - private IProducer GetProducer() + private IProducer GetProducer() { if (_producer != null) return _producer; @@ -79,7 +80,9 @@ private IProducer GetProducer() ClientId = _configuration.ClientId }; - _producer = new ProducerBuilder(producerConfig).Build(); + _configuration.ApplyBrokerAuthentication(producerConfig); + + _producer = new ProducerBuilder(producerConfig).Build(); return _producer; } } @@ -93,7 +96,7 @@ private static Headers BuildHeaders(IMessage message) foreach (var userProperty in message.UserProperties) { - if (userProperty.Value == null) + if (userProperty.Value is null) continue; headers.Add(userProperty.Key, Encoding.UTF8.GetBytes(userProperty.Value.ToString()!)); @@ -106,6 +109,7 @@ public void Dispose() { _producer?.Flush(TimeSpan.FromSeconds(10)); _producer?.Dispose(); + _kafkaMessageSerializer.DisposeAsync().AsTask().GetAwaiter().GetResult(); GC.SuppressFinalize(this); } } From e6a44473f1a7af29d6c3608fd805a0bc4735b573 Mon Sep 17 00:00:00 2001 From: Fabio Della Rosa Date: Thu, 18 Jun 2026 10:55:10 +0200 Subject: [PATCH 2/4] Add portions of docker-compose to manage kafka implementation --- .gitignore | 1 + docker-samples/docker-compose.sasl.yml | 78 ++++++++++++++++++++++++++ docker-samples/docker-compose.yml | 62 ++++++++++++++++++++ docker-samples/kafka_server_jaas.conf | 12 ++++ 4 files changed, 153 insertions(+) create mode 100644 docker-samples/docker-compose.sasl.yml create mode 100644 docker-samples/docker-compose.yml create mode 100644 docker-samples/kafka_server_jaas.conf diff --git a/.gitignore b/.gitignore index ce89292..f14fc74 100644 --- a/.gitignore +++ b/.gitignore @@ -416,3 +416,4 @@ FodyWeavers.xsd *.msix *.msm *.msp +/.dotnet diff --git a/docker-samples/docker-compose.sasl.yml b/docker-samples/docker-compose.sasl.yml new file mode 100644 index 0000000..b36c6ac --- /dev/null +++ b/docker-samples/docker-compose.sasl.yml @@ -0,0 +1,78 @@ + +# This docker-compose file sets up a Kafka cluster with authentication, along with a schema registry and a Kafka UI for management +# Required the kafka_server_jaas.conf file to be present in the same directory. +services: + kafka: + image: confluentinc/cp-kafka:latest + container_name: kafka + ports: + - "9092:9092" + networks: + - my-network + volumes: + - ./kafka_server_jaas.conf:/etc/kafka/kafka_server_jaas.conf:ro + environment: + KAFKA_OPTS: "-Djava.security.auth.login.config=/etc/kafka/kafka_server_jaas.conf" + + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + + KAFKA_LISTENERS: INTERNAL://0.0.0.0:29092,EXTERNAL://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 + KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:29092,EXTERNAL://localhost:9092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:SASL_PLAINTEXT,EXTERNAL:SASL_PLAINTEXT,CONTROLLER:PLAINTEXT + + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093 + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL + + KAFKA_SASL_ENABLED_MECHANISMS: PLAIN + KAFKA_SASL_MECHANISM_INTER_BROKER_PROTOCOL: PLAIN + KAFKA_LISTENER_NAME_INTERNAL_PLAIN_SASL_JAAS_CONFIG: 'org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="admin-secret" user_admin="admin-secret";' + KAFKA_LISTENER_NAME_EXTERNAL_PLAIN_SASL_JAAS_CONFIG: 'org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="admin-secret" user_admin="admin-secret";' + + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + + CLUSTER_ID: + + schema-registry: + image: confluentinc/cp-schema-registry:latest + container_name: schema-registry + depends_on: + - kafka + ports: + - "8081:8081" + networks: + - my-network + environment: + SCHEMA_REGISTRY_HOST_NAME: schema-registry + SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: SASL_PLAINTEXT://kafka:29092 + SCHEMA_REGISTRY_KAFKASTORE_SECURITY_PROTOCOL: SASL_PLAINTEXT + SCHEMA_REGISTRY_KAFKASTORE_SASL_MECHANISM: PLAIN + SCHEMA_REGISTRY_KAFKASTORE_SASL_JAAS_CONFIG: 'org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="admin-secret";' + SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 + + kafka-ui: + image: ghcr.io/kafbat/kafka-ui:latest + container_name: kafka-ui + depends_on: + - kafka + - schema-registry + ports: + - "7080:8080" + networks: + - my-network + environment: + KAFKA_CLUSTERS_0_NAME: kafka0 + KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092 + KAFKA_CLUSTERS_0_SCHEMAREGISTRY: http://schema-registry:8081 + KAFKA_CLUSTERS_0_PROPERTIES_SECURITY_PROTOCOL: SASL_PLAINTEXT + KAFKA_CLUSTERS_0_PROPERTIES_SASL_MECHANISM: PLAIN + KAFKA_CLUSTERS_0_PROPERTIES_SASL_JAAS_CONFIG: 'org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="admin-secret";' + DYNAMIC_CONFIG_ENABLED: "true" + +networks: + my-network: + driver: bridge \ No newline at end of file diff --git a/docker-samples/docker-compose.yml b/docker-samples/docker-compose.yml new file mode 100644 index 0000000..034f794 --- /dev/null +++ b/docker-samples/docker-compose.yml @@ -0,0 +1,62 @@ +# This docker-compose file use kafka in KRaft mode without zookeeper, no authentication, and no authorization. It is intended for local development and testing purposes only. + +services: + kafka: + image: confluentinc/cp-kafka:latest + container_name: kafka + ports: + - "9092:9092" + networks: + - my-network + environment: + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + + KAFKA_LISTENERS: INTERNAL://0.0.0.0:29092,EXTERNAL://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 + KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:29092,EXTERNAL://localhost:9092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT + + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093 + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL + + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" + + CLUSTER_ID: + + schema-registry: + image: confluentinc/cp-schema-registry:latest + container_name: schema-registry + depends_on: + - kafka + ports: + - "8081:8081" + networks: + - my-network + environment: + SCHEMA_REGISTRY_HOST_NAME: schema-registry + SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka:29092 + SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 + + kafka-ui: + image: ghcr.io/kafbat/kafka-ui:latest + container_name: kafka-ui + depends_on: + - kafka + - schema-registry + ports: + - "7080:8080" + networks: + - my-network + environment: + KAFKA_CLUSTERS_0_NAME: kafka0 + KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092 + KAFKA_CLUSTERS_0_SCHEMAREGISTRY: http://schema-registry:8081 + DYNAMIC_CONFIG_ENABLED: "true" + +networks: + my-network: + driver: bridge \ No newline at end of file diff --git a/docker-samples/kafka_server_jaas.conf b/docker-samples/kafka_server_jaas.conf new file mode 100644 index 0000000..ba70ffc --- /dev/null +++ b/docker-samples/kafka_server_jaas.conf @@ -0,0 +1,12 @@ +KafkaServer { + org.apache.kafka.common.security.plain.PlainLoginModule required + username="admin" + password="admin-secret" + user_admin="admin-secret"; +}; + +Client { + org.apache.kafka.common.security.plain.PlainLoginModule required + username="admin" + password="admin-secret"; +}; \ No newline at end of file From 328297de0e2bb6af056d46df3def74355fb4d2ec Mon Sep 17 00:00:00 2001 From: "fabio.dellarosa" Date: Mon, 22 Jun 2026 10:08:36 +0200 Subject: [PATCH 3/4] Enhance README with comprehensive project documentation Added detailed documentation for Muflone.Transport.Kafka, including breaking changes, features, installation instructions, architecture overview, quick start guide, configuration options, and registered services. --- README.md | 340 +++++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 339 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 1b74df6..63df329 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,340 @@ # Muflone.Transport.Kafka -Package to integrate Kafka with Muflone + +[![NuGet](https://img.shields.io/nuget/v/Muflone.Transport.Kafka.svg)](https://www.nuget.org/packages/Muflone.Transport.Kafka/) +[![License: MIT](https://img.shields.io/badge/License-MIT-yellow.svg)](https://opensource.org/licenses/MIT) + +Muflone extension to manage messages on Kafka, designed for CQRS and Event Sourcing architectures. + +## Breaking Changes + +### v10.1.0 - Handler subscription alignment + +Kafka now supports the same `IMessageSubscriber` and `MessageHandlersStarter` integration style used by `Muflone.Transport.RabbitMQ`. + +Explicit consumers are still supported, but the transport can now also dispatch handlers through Muflone's subscriber infrastructure without relying on the old consumer-subscription approach. + +--- + +## Features + +- **CQRS Support** - Commands and events are transported through Kafka topics with consumer-group based routing +- **Kafka-native Messaging** - Producers publish serialized Muflone messages with `UserProperties` preserved as Kafka headers +- **Automatic Handler Dispatch** - Supports `IMessageSubscriber` and `MessageHandlersStarter`, aligned with the RabbitMQ transport +- **Explicit Consumers** - Supports `CommandConsumerBase`, `DomainEventsConsumerBase`, and `IntegrationEventsConsumerBase` +- **Topic Naming Convention** - Customize topic naming through `KafkaConfiguration.TopicNamingConvention` +- **Schema Registry Integration** - Supports Confluent Schema Registry for `Json`, `Avro`, and `Protobuf` +- **Basic Auth** - Supports SASL basic auth for Kafka brokers and HTTP Basic Auth for Schema Registry +- **Async-first** - All operations use async/await +- **.NET 10** - Built for `net10.0` with nullable reference types + +## Install + +``` +dotnet add package Muflone.Transport.Kafka +``` + +or + +``` +Install-Package Muflone.Transport.Kafka +``` + +## Architecture Overview + +Kafka does not differentiate between commands and events at the broker level, so the library models both with topics: + +| Concept | Transport | Routing | Consumer Behavior | +|--------------------|-------------|----------------------------------------------|-----------------------------------------| +| Commands | Kafka topic | One logical consumer per consumer group | Competing consumers | +| Domain Events | Kafka topic | One-to-many across different consumer groups | Pub/sub through separate groups | +| Integration Events | Kafka topic | One-to-many across different consumer groups | Pub/sub through separate groups | + +**Commands** are sent to a topic and processed by one consumer inside the same consumer group. +**Domain Events** and **Integration Events** are published to topics and can be consumed by multiple subscribers by using different consumer groups. + +By default, each message type uses a topic named after the CLR type in lowercase: + +- `CreateOrder` -> `createorder` +- `OrderCreated` -> `ordercreated` + +## Quick Start + +### 1. Define your messages + +Commands and events must extend the Muflone base classes: + +```csharp +public class CreateOrder : Command +{ + public readonly string OrderNumber; + + public CreateOrder(OrderId aggregateId, string orderNumber) + : base(aggregateId) + { + OrderNumber = orderNumber; + } +} + +public class OrderCreated : DomainEvent +{ + public readonly string OrderNumber; + + public OrderCreated(OrderId aggregateId, string orderNumber) + : base(aggregateId) + { + OrderNumber = orderNumber; + } +} +``` + +### 2. Create command and event handlers + +```csharp +public class CreateOrderCommandHandler : ICommandHandlerAsync +{ + private readonly IRepository _repository; + + public CreateOrderCommandHandler(IRepository repository) + { + _repository = repository; + } + + public async Task HandleAsync(CreateOrder command, CancellationToken cancellationToken = default) + { + var order = Order.CreateOrder(command.AggregateId, command.OrderNumber); + await _repository.SaveAsync(order, Guid.NewGuid(), cancellationToken); + } +} + +public class OrderCreatedEventHandler : IDomainEventHandlerAsync +{ + private readonly ILogger _logger; + + public OrderCreatedEventHandler(ILoggerFactory loggerFactory) + { + _logger = loggerFactory.CreateLogger(GetType()); + } + + public Task HandleAsync(OrderCreated @event, CancellationToken cancellationToken = default) + { + _logger.LogInformation("Order created: {OrderNumber}", @event.OrderNumber); + return Task.CompletedTask; + } +} +``` + +### 3. Create consumers + +Consumers wire messages to their handlers. Extend the appropriate base class: + +```csharp +// Command consumer (one handler per command) +public class CreateOrderConsumer : CommandConsumerBase +{ + protected override ICommandHandlerAsync HandlerAsync { get; } + + public CreateOrderConsumer( + IRepository repository, + KafkaConfiguration configuration, + ILoggerFactory loggerFactory) + : base(repository, configuration, loggerFactory) + { + HandlerAsync = new CreateOrderCommandHandler(repository); + } +} + +// Domain event consumer (supports multiple handlers per event) +public class OrderCreatedConsumer : DomainEventsConsumerBase +{ + protected override IEnumerable> HandlersAsync { get; } + + public OrderCreatedConsumer( + KafkaConfiguration configuration, + ILoggerFactory loggerFactory) + : base(configuration, loggerFactory) + { + HandlersAsync = new List> + { + new OrderCreatedEventHandler(loggerFactory) + }; + } +} + +// Integration event consumer (for cross-boundary events) +public class OrderShippedConsumer : IntegrationEventsConsumerBase +{ + protected override IEnumerable> HandlersAsync { get; } + + public OrderShippedConsumer( + KafkaConfiguration configuration, + ILoggerFactory loggerFactory) + : base(configuration, loggerFactory) + { + HandlersAsync = new List> + { + new OrderShippedIntegrationHandler(loggerFactory) + }; + } +} +``` + +### 4. Register services in DI + +```csharp +var builder = WebApplication.CreateBuilder(args); + +var loggerFactory = LoggerFactory.Create(logging => logging.AddConsole()); + +// Configure Kafka connection +var kafkaConfiguration = new KafkaConfiguration( + bootstrapServers: "localhost:9092", + groupId: "OrderService", + clientId: "OrderService", + autoOffsetReset: AutoOffsetReset.Earliest +); + +kafkaConfiguration.BrokerUsername = "kafka-user"; +kafkaConfiguration.BrokerPassword = "kafka-password"; +kafkaConfiguration.SchemaRegistryUrl = "https://localhost:8081"; +kafkaConfiguration.SchemaRegistryUsername = "schema-user"; +kafkaConfiguration.SchemaRegistryPassword = "schema-password"; +kafkaConfiguration.SerializationFormat = KafkaSerializationFormat.Json; + +// Register handlers +builder.Services.AddScoped, CreateOrderCommandHandler>(); +builder.Services.AddScoped, OrderCreatedEventHandler>(); + +// Register Muflone Kafka transport +builder.Services.AddMufloneTransportKafka(loggerFactory, kafkaConfiguration); +``` + +If you want to register explicit consumers too: + +```csharp +builder.Services.AddMufloneKafkaConsumers( + new IConsumer[] + { + new CreateOrderConsumer(repository, kafkaConfiguration, loggerFactory), + new OrderCreatedConsumer(kafkaConfiguration, loggerFactory) + }); +``` + +### 5. Send commands and publish events + +```csharp +// Inject IServiceBus to send commands +public class OrderController : ControllerBase +{ + private readonly IServiceBus _serviceBus; + + public OrderController(IServiceBus serviceBus) + { + _serviceBus = serviceBus; + } + + [HttpPost] + public async Task CreateOrder(CreateOrderRequest request) + { + var command = new CreateOrder( + new OrderId(Guid.NewGuid()), + request.OrderNumber); + + await _serviceBus.SendAsync(command); + return Accepted(); + } +} + +// Inject IEventBus to publish events +public class OrderService +{ + private readonly IEventBus _eventBus; + + public OrderService(IEventBus eventBus) + { + _eventBus = eventBus; + } + + public async Task NotifyOrderCreated(OrderId orderId, string orderNumber) + { + var @event = new OrderCreated(orderId, orderNumber); + await _eventBus.PublishAsync(@event); + } +} +``` + +## Configuration + +`KafkaConfiguration` supports several constructor overloads: + +```csharp +// Basic (default clientId = groupId, default AutoOffsetReset = Earliest) +new KafkaConfiguration(bootstrapServers, groupId); + +// With explicit offset reset +new KafkaConfiguration(bootstrapServers, groupId, AutoOffsetReset.Earliest); + +// With explicit clientId and offset reset +new KafkaConfiguration(bootstrapServers, groupId, clientId, AutoOffsetReset.Earliest); +``` + +| Parameter | Description | Default | +|---------------------|-----------------------------------------------------|------------| +| `bootstrapServers` | Kafka broker list, for example `localhost:9092` | - | +| `groupId` | Consumer group id used by Kafka consumers | - | +| `clientId` | Client id used by the producer and subscribers | `groupId` | +| `autoOffsetReset` | Kafka offset reset strategy | `Earliest` | + +### Broker basic auth + +```csharp +kafkaConfiguration.BrokerUsername = "kafka-user"; +kafkaConfiguration.BrokerPassword = "kafka-password"; +kafkaConfiguration.BrokerSecurityProtocol = SecurityProtocol.SaslSsl; +kafkaConfiguration.BrokerSaslMechanism = SaslMechanism.Plain; +``` + +### Schema Registry and serialization format + +The transport keeps using the configured Muflone `ISerializer` for the domain message payload and wraps that payload in a Kafka envelope managed through Schema Registry. + +This means: + +- existing Muflone messages do not need to become Avro or Protobuf classes +- Kafka producer/consumer wire format is controlled by `KafkaSerializationFormat` +- Schema Registry subjects are registered through the Confluent .NET serdes packages + +```csharp +kafkaConfiguration.SchemaRegistryUrl = "https://localhost:8081"; +kafkaConfiguration.SchemaRegistryUsername = "schema-user"; +kafkaConfiguration.SchemaRegistryPassword = "schema-password"; +kafkaConfiguration.SchemaRegistryAutoRegisterSchemas = true; +kafkaConfiguration.SerializationFormat = KafkaSerializationFormat.Avro; +// or KafkaSerializationFormat.Json +// or KafkaSerializationFormat.Protobuf +``` + + +### Custom topic naming + +```csharp +kafkaConfiguration.TopicNamingConvention = type => $"{type.Name.ToLowerInvariant()}"; +``` + +## Registered Services + +`AddMufloneTransportKafka` registers the following services as singletons: + +| Interface | Implementation | +|----------------------|------------------------| +| `IServiceBus` | `ServiceBus` | +| `IEventBus` | `ServiceBus` | +| `IMessageSubscriber` | `KafkaSubscriber` | +| `IHostedService` | `KafkaBrokerStarter` | +| `IHostedService` | `MessageHandlersStarter` | + +## Fully working example +TODO + +## License + +This project is licensed under the [MIT License](https://opensource.org/licenses/MIT). From d87b874f9c36b8641b0d9df6453c572f61c61408 Mon Sep 17 00:00:00 2001 From: "fabio.dellarosa" Date: Mon, 22 Jun 2026 10:08:36 +0200 Subject: [PATCH 4/4] Enhance README with comprehensive project documentation Added detailed documentation for Muflone.Transport.Kafka, including breaking changes, features, installation instructions, architecture overview, quick start guide, configuration options, and registered services. --- README.md | 340 +++++++++++++++++- .../Muflone.Transport.Kafka.csproj | 6 + 2 files changed, 345 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 1b74df6..63df329 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,340 @@ # Muflone.Transport.Kafka -Package to integrate Kafka with Muflone + +[![NuGet](https://img.shields.io/nuget/v/Muflone.Transport.Kafka.svg)](https://www.nuget.org/packages/Muflone.Transport.Kafka/) +[![License: MIT](https://img.shields.io/badge/License-MIT-yellow.svg)](https://opensource.org/licenses/MIT) + +Muflone extension to manage messages on Kafka, designed for CQRS and Event Sourcing architectures. + +## Breaking Changes + +### v10.1.0 - Handler subscription alignment + +Kafka now supports the same `IMessageSubscriber` and `MessageHandlersStarter` integration style used by `Muflone.Transport.RabbitMQ`. + +Explicit consumers are still supported, but the transport can now also dispatch handlers through Muflone's subscriber infrastructure without relying on the old consumer-subscription approach. + +--- + +## Features + +- **CQRS Support** - Commands and events are transported through Kafka topics with consumer-group based routing +- **Kafka-native Messaging** - Producers publish serialized Muflone messages with `UserProperties` preserved as Kafka headers +- **Automatic Handler Dispatch** - Supports `IMessageSubscriber` and `MessageHandlersStarter`, aligned with the RabbitMQ transport +- **Explicit Consumers** - Supports `CommandConsumerBase`, `DomainEventsConsumerBase`, and `IntegrationEventsConsumerBase` +- **Topic Naming Convention** - Customize topic naming through `KafkaConfiguration.TopicNamingConvention` +- **Schema Registry Integration** - Supports Confluent Schema Registry for `Json`, `Avro`, and `Protobuf` +- **Basic Auth** - Supports SASL basic auth for Kafka brokers and HTTP Basic Auth for Schema Registry +- **Async-first** - All operations use async/await +- **.NET 10** - Built for `net10.0` with nullable reference types + +## Install + +``` +dotnet add package Muflone.Transport.Kafka +``` + +or + +``` +Install-Package Muflone.Transport.Kafka +``` + +## Architecture Overview + +Kafka does not differentiate between commands and events at the broker level, so the library models both with topics: + +| Concept | Transport | Routing | Consumer Behavior | +|--------------------|-------------|----------------------------------------------|-----------------------------------------| +| Commands | Kafka topic | One logical consumer per consumer group | Competing consumers | +| Domain Events | Kafka topic | One-to-many across different consumer groups | Pub/sub through separate groups | +| Integration Events | Kafka topic | One-to-many across different consumer groups | Pub/sub through separate groups | + +**Commands** are sent to a topic and processed by one consumer inside the same consumer group. +**Domain Events** and **Integration Events** are published to topics and can be consumed by multiple subscribers by using different consumer groups. + +By default, each message type uses a topic named after the CLR type in lowercase: + +- `CreateOrder` -> `createorder` +- `OrderCreated` -> `ordercreated` + +## Quick Start + +### 1. Define your messages + +Commands and events must extend the Muflone base classes: + +```csharp +public class CreateOrder : Command +{ + public readonly string OrderNumber; + + public CreateOrder(OrderId aggregateId, string orderNumber) + : base(aggregateId) + { + OrderNumber = orderNumber; + } +} + +public class OrderCreated : DomainEvent +{ + public readonly string OrderNumber; + + public OrderCreated(OrderId aggregateId, string orderNumber) + : base(aggregateId) + { + OrderNumber = orderNumber; + } +} +``` + +### 2. Create command and event handlers + +```csharp +public class CreateOrderCommandHandler : ICommandHandlerAsync +{ + private readonly IRepository _repository; + + public CreateOrderCommandHandler(IRepository repository) + { + _repository = repository; + } + + public async Task HandleAsync(CreateOrder command, CancellationToken cancellationToken = default) + { + var order = Order.CreateOrder(command.AggregateId, command.OrderNumber); + await _repository.SaveAsync(order, Guid.NewGuid(), cancellationToken); + } +} + +public class OrderCreatedEventHandler : IDomainEventHandlerAsync +{ + private readonly ILogger _logger; + + public OrderCreatedEventHandler(ILoggerFactory loggerFactory) + { + _logger = loggerFactory.CreateLogger(GetType()); + } + + public Task HandleAsync(OrderCreated @event, CancellationToken cancellationToken = default) + { + _logger.LogInformation("Order created: {OrderNumber}", @event.OrderNumber); + return Task.CompletedTask; + } +} +``` + +### 3. Create consumers + +Consumers wire messages to their handlers. Extend the appropriate base class: + +```csharp +// Command consumer (one handler per command) +public class CreateOrderConsumer : CommandConsumerBase +{ + protected override ICommandHandlerAsync HandlerAsync { get; } + + public CreateOrderConsumer( + IRepository repository, + KafkaConfiguration configuration, + ILoggerFactory loggerFactory) + : base(repository, configuration, loggerFactory) + { + HandlerAsync = new CreateOrderCommandHandler(repository); + } +} + +// Domain event consumer (supports multiple handlers per event) +public class OrderCreatedConsumer : DomainEventsConsumerBase +{ + protected override IEnumerable> HandlersAsync { get; } + + public OrderCreatedConsumer( + KafkaConfiguration configuration, + ILoggerFactory loggerFactory) + : base(configuration, loggerFactory) + { + HandlersAsync = new List> + { + new OrderCreatedEventHandler(loggerFactory) + }; + } +} + +// Integration event consumer (for cross-boundary events) +public class OrderShippedConsumer : IntegrationEventsConsumerBase +{ + protected override IEnumerable> HandlersAsync { get; } + + public OrderShippedConsumer( + KafkaConfiguration configuration, + ILoggerFactory loggerFactory) + : base(configuration, loggerFactory) + { + HandlersAsync = new List> + { + new OrderShippedIntegrationHandler(loggerFactory) + }; + } +} +``` + +### 4. Register services in DI + +```csharp +var builder = WebApplication.CreateBuilder(args); + +var loggerFactory = LoggerFactory.Create(logging => logging.AddConsole()); + +// Configure Kafka connection +var kafkaConfiguration = new KafkaConfiguration( + bootstrapServers: "localhost:9092", + groupId: "OrderService", + clientId: "OrderService", + autoOffsetReset: AutoOffsetReset.Earliest +); + +kafkaConfiguration.BrokerUsername = "kafka-user"; +kafkaConfiguration.BrokerPassword = "kafka-password"; +kafkaConfiguration.SchemaRegistryUrl = "https://localhost:8081"; +kafkaConfiguration.SchemaRegistryUsername = "schema-user"; +kafkaConfiguration.SchemaRegistryPassword = "schema-password"; +kafkaConfiguration.SerializationFormat = KafkaSerializationFormat.Json; + +// Register handlers +builder.Services.AddScoped, CreateOrderCommandHandler>(); +builder.Services.AddScoped, OrderCreatedEventHandler>(); + +// Register Muflone Kafka transport +builder.Services.AddMufloneTransportKafka(loggerFactory, kafkaConfiguration); +``` + +If you want to register explicit consumers too: + +```csharp +builder.Services.AddMufloneKafkaConsumers( + new IConsumer[] + { + new CreateOrderConsumer(repository, kafkaConfiguration, loggerFactory), + new OrderCreatedConsumer(kafkaConfiguration, loggerFactory) + }); +``` + +### 5. Send commands and publish events + +```csharp +// Inject IServiceBus to send commands +public class OrderController : ControllerBase +{ + private readonly IServiceBus _serviceBus; + + public OrderController(IServiceBus serviceBus) + { + _serviceBus = serviceBus; + } + + [HttpPost] + public async Task CreateOrder(CreateOrderRequest request) + { + var command = new CreateOrder( + new OrderId(Guid.NewGuid()), + request.OrderNumber); + + await _serviceBus.SendAsync(command); + return Accepted(); + } +} + +// Inject IEventBus to publish events +public class OrderService +{ + private readonly IEventBus _eventBus; + + public OrderService(IEventBus eventBus) + { + _eventBus = eventBus; + } + + public async Task NotifyOrderCreated(OrderId orderId, string orderNumber) + { + var @event = new OrderCreated(orderId, orderNumber); + await _eventBus.PublishAsync(@event); + } +} +``` + +## Configuration + +`KafkaConfiguration` supports several constructor overloads: + +```csharp +// Basic (default clientId = groupId, default AutoOffsetReset = Earliest) +new KafkaConfiguration(bootstrapServers, groupId); + +// With explicit offset reset +new KafkaConfiguration(bootstrapServers, groupId, AutoOffsetReset.Earliest); + +// With explicit clientId and offset reset +new KafkaConfiguration(bootstrapServers, groupId, clientId, AutoOffsetReset.Earliest); +``` + +| Parameter | Description | Default | +|---------------------|-----------------------------------------------------|------------| +| `bootstrapServers` | Kafka broker list, for example `localhost:9092` | - | +| `groupId` | Consumer group id used by Kafka consumers | - | +| `clientId` | Client id used by the producer and subscribers | `groupId` | +| `autoOffsetReset` | Kafka offset reset strategy | `Earliest` | + +### Broker basic auth + +```csharp +kafkaConfiguration.BrokerUsername = "kafka-user"; +kafkaConfiguration.BrokerPassword = "kafka-password"; +kafkaConfiguration.BrokerSecurityProtocol = SecurityProtocol.SaslSsl; +kafkaConfiguration.BrokerSaslMechanism = SaslMechanism.Plain; +``` + +### Schema Registry and serialization format + +The transport keeps using the configured Muflone `ISerializer` for the domain message payload and wraps that payload in a Kafka envelope managed through Schema Registry. + +This means: + +- existing Muflone messages do not need to become Avro or Protobuf classes +- Kafka producer/consumer wire format is controlled by `KafkaSerializationFormat` +- Schema Registry subjects are registered through the Confluent .NET serdes packages + +```csharp +kafkaConfiguration.SchemaRegistryUrl = "https://localhost:8081"; +kafkaConfiguration.SchemaRegistryUsername = "schema-user"; +kafkaConfiguration.SchemaRegistryPassword = "schema-password"; +kafkaConfiguration.SchemaRegistryAutoRegisterSchemas = true; +kafkaConfiguration.SerializationFormat = KafkaSerializationFormat.Avro; +// or KafkaSerializationFormat.Json +// or KafkaSerializationFormat.Protobuf +``` + + +### Custom topic naming + +```csharp +kafkaConfiguration.TopicNamingConvention = type => $"{type.Name.ToLowerInvariant()}"; +``` + +## Registered Services + +`AddMufloneTransportKafka` registers the following services as singletons: + +| Interface | Implementation | +|----------------------|------------------------| +| `IServiceBus` | `ServiceBus` | +| `IEventBus` | `ServiceBus` | +| `IMessageSubscriber` | `KafkaSubscriber` | +| `IHostedService` | `KafkaBrokerStarter` | +| `IHostedService` | `MessageHandlersStarter` | + +## Fully working example +TODO + +## License + +This project is licensed under the [MIT License](https://opensource.org/licenses/MIT). diff --git a/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj b/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj index 2d0fd4e..a35b17e 100644 --- a/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj +++ b/src/Muflone.Transport.Kafka/Muflone.Transport.Kafka.csproj @@ -53,4 +53,10 @@ + + + Always + + +