From 4651b7a46d2fb66d449c8cad36ead9d95bc1fd2a Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 7 Dec 2024 20:07:30 +0100 Subject: [PATCH] Add schema registry support to Kafka module (#6190) This commit introduces a schema registry functionality by adding interfaces and classes to manage schema registry definitions in the Elsa.Kafka module. It updates KafkaOptions to include schema registries and modifies classes to support schema registry configurations for producers and consumers. Additionally, it updates package references to include necessary dependencies for schema registry support. --- Directory.Packages.props | 3 ++- src/apps/Elsa.Server.Web/appsettings.json | 15 ++++++++++-- .../Elsa.Kafka/Activities/ProduceMessage.cs | 20 ++++++++++++---- .../ISchemaRegistryDefinitionEnumerator.cs | 11 +++++++++ .../ISchemaRegistryDefinitionProvider.cs | 6 +++++ src/modules/Elsa.Kafka/Elsa.Kafka.csproj | 1 + .../Elsa.Kafka/Entities/ConsumerDefinition.cs | 2 +- .../Elsa.Kafka/Entities/ProducerDefinition.cs | 2 +- .../Entities/SchemaRegistryDefinition.cs | 10 ++++++++ .../Elsa.Kafka/Features/KafkaFeature.cs | 9 +++++--- .../SchemaRegistryDefinitionEnumerator.cs | 23 +++++++++++++++++++ .../Implementations/WorkerManager.cs | 19 +++++++++++---- .../Models/CreateConsumerContext.cs | 2 +- .../Models/CreateProducerContext.cs | 2 +- .../Elsa.Kafka/Options/KafkaOptions.cs | 1 + .../Providers/OptionsDefinitionProvider.cs | 7 +++++- 16 files changed, 113 insertions(+), 20 deletions(-) create mode 100644 src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionEnumerator.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionProvider.cs create mode 100644 src/modules/Elsa.Kafka/Entities/SchemaRegistryDefinition.cs create mode 100644 src/modules/Elsa.Kafka/Implementations/SchemaRegistryDefinitionEnumerator.cs diff --git a/Directory.Packages.props b/Directory.Packages.props index 157742de1..4bdcfbbce 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -15,7 +15,8 @@ - + + diff --git a/src/apps/Elsa.Server.Web/appsettings.json b/src/apps/Elsa.Server.Web/appsettings.json index ea33ad3b9..1f6fd76cd 100644 --- a/src/apps/Elsa.Server.Web/appsettings.json +++ b/src/apps/Elsa.Server.Web/appsettings.json @@ -217,6 +217,15 @@ "Name": "topic-2" } ], + "SchemaRegistries": [ + { + "Id": "schema-registry-1", + "Name": "Schema Registry 1", + "Config": { + "Url": "http://localhost:8081" + } + } + ], "Producers": [ { "Id": "producer-1", @@ -224,7 +233,8 @@ "FactoryType": "Elsa.Kafka.Factories.GenericProducerFactory`2[[System.String, System.Private.CoreLib], [Elsa.Server.Web.Messages.OrderReceived, Elsa.Server.Web]], Elsa.Kafka", "Config": { "BootstrapServers": "localhost:9092" - } + }, + "SchemaRegistryId": "schema-registry-1" } ], "Consumers": [ @@ -237,7 +247,8 @@ "GroupId": "group-1", "AutoOffsetReset": "earliest", "EnableAutoCommit": "true" - } + }, + "SchemaRegistryId": "schema-registry-1" } ] }, diff --git a/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs b/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs index 01d92bceb..a4536eb7d 100644 --- a/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs +++ b/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs @@ -5,10 +5,9 @@ using Elsa.Kafka.UIHints; using Elsa.Workflows; using Elsa.Workflows.Attributes; using Elsa.Workflows.Models; -using Elsa.Workflows.Runtime; using Elsa.Workflows.UIHints; -using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; +using Open.Linq.AsyncExtensions; namespace Elsa.Kafka.Activities; @@ -73,7 +72,7 @@ public class ProduceMessage : CodeActivity if (key is string keyString && string.IsNullOrWhiteSpace(keyString)) key = null; - using var producer = CreateProducer(context, producerDefinition); + using var producer = await CreateProducerAsync(context, producerDefinition); var headers = CreateHeaders(context); await producer.ProduceAsync(topic, key, content, headers, cancellationToken); } @@ -94,14 +93,25 @@ public class ProduceMessage : CodeActivity return headers; } - private IProducer CreateProducer(ActivityExecutionContext context, ProducerDefinition producerDefinition) + private async Task CreateProducerAsync(ActivityExecutionContext context, ProducerDefinition producerDefinition) { var factory = context.GetOrCreateService(producerDefinition.FactoryType) as IProducerFactory; if (factory == null) throw new InvalidOperationException($"Producer factory of type '{producerDefinition.FactoryType}' not found."); - var createProducerContext = new CreateProducerContext(producerDefinition); + var schemaRegistryDefinition = await GetSchemaRegistryDefinitionAsync(context, producerDefinition.SchemaRegistryId); + var createProducerContext = new CreateProducerContext(producerDefinition, schemaRegistryDefinition); return factory.CreateProducer(createProducerContext); } + + private async Task GetSchemaRegistryDefinitionAsync(ActivityExecutionContext context, string? id, CancellationToken cancellationToken = default) + { + if (id == null) + return null; + + var schemaRegistryDefinitionEnumerator = context.GetRequiredService(); + var schemaRegistryDefinitions = await schemaRegistryDefinitionEnumerator.EnumerateAsync(cancellationToken).ToList(); + return schemaRegistryDefinitions.FirstOrDefault(x => x.Id == id); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionEnumerator.cs new file mode 100644 index 000000000..e0240a432 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionEnumerator.cs @@ -0,0 +1,11 @@ +namespace Elsa.Kafka; + +public interface ISchemaRegistryDefinitionEnumerator +{ + /// + /// Retrieves a list of schema registry definitions provided by implementations. + /// + /// A token to monitor for cancellation requests. + /// A list of SchemaRegistryDefinition objects. + Task> EnumerateAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionProvider.cs b/src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionProvider.cs new file mode 100644 index 000000000..cb38bee71 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/ISchemaRegistryDefinitionProvider.cs @@ -0,0 +1,6 @@ +namespace Elsa.Kafka; + +public interface ISchemaRegistryDefinitionProvider +{ + Task> GetSchemaRegistryDefinitionsAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Elsa.Kafka.csproj b/src/modules/Elsa.Kafka/Elsa.Kafka.csproj index fb8e823ad..0f9959597 100644 --- a/src/modules/Elsa.Kafka/Elsa.Kafka.csproj +++ b/src/modules/Elsa.Kafka/Elsa.Kafka.csproj @@ -11,6 +11,7 @@ + diff --git a/src/modules/Elsa.Kafka/Entities/ConsumerDefinition.cs b/src/modules/Elsa.Kafka/Entities/ConsumerDefinition.cs index b1d9cb124..f9320d326 100644 --- a/src/modules/Elsa.Kafka/Entities/ConsumerDefinition.cs +++ b/src/modules/Elsa.Kafka/Entities/ConsumerDefinition.cs @@ -1,6 +1,5 @@ using Confluent.Kafka; using Elsa.Common.Entities; -using Elsa.Kafka.Factories; namespace Elsa.Kafka; @@ -9,4 +8,5 @@ public class ConsumerDefinition : Entity public string Name { get; set; } = default!; public Type FactoryType { get; set; } = default!; public ConsumerConfig Config { get; set; } = new(); + public string? SchemaRegistryId { get; set; } = default!; } \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs b/src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs index 1544aa243..b2f410329 100644 --- a/src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs +++ b/src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs @@ -1,6 +1,5 @@ using Confluent.Kafka; using Elsa.Common.Entities; -using Elsa.Kafka.Factories; namespace Elsa.Kafka; @@ -9,4 +8,5 @@ public class ProducerDefinition : Entity public string Name { get; set; } = default!; public Type FactoryType { get; set; } = default!; public ProducerConfig Config { get; set; } = new(); + public string? SchemaRegistryId { get; set; } = default!; } \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Entities/SchemaRegistryDefinition.cs b/src/modules/Elsa.Kafka/Entities/SchemaRegistryDefinition.cs new file mode 100644 index 000000000..0b8620867 --- /dev/null +++ b/src/modules/Elsa.Kafka/Entities/SchemaRegistryDefinition.cs @@ -0,0 +1,10 @@ +using Confluent.SchemaRegistry; +using Elsa.Common.Entities; + +namespace Elsa.Kafka; + +public class SchemaRegistryDefinition : Entity +{ + public string Name { get; set; } = default!; + public SchemaRegistryConfig Config { get; set; } = default!; +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Features/KafkaFeature.cs b/src/modules/Elsa.Kafka/Features/KafkaFeature.cs index 35a41a980..dd4067c83 100644 --- a/src/modules/Elsa.Kafka/Features/KafkaFeature.cs +++ b/src/modules/Elsa.Kafka/Features/KafkaFeature.cs @@ -60,12 +60,15 @@ public class KafkaFeature(IModule module) : FeatureBase(module) .AddBackgroundTask() .AddSingleton() .AddScoped() - .AddScoped() - .AddScoped() - .AddScoped() + .AddScoped() + .AddScoped(sp => sp.GetRequiredService()) + .AddScoped(sp => sp.GetRequiredService()) + .AddScoped(sp => sp.GetRequiredService()) + .AddScoped(sp => sp.GetRequiredService()) .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Kafka/Implementations/SchemaRegistryDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Implementations/SchemaRegistryDefinitionEnumerator.cs new file mode 100644 index 000000000..30e635cc5 --- /dev/null +++ b/src/modules/Elsa.Kafka/Implementations/SchemaRegistryDefinitionEnumerator.cs @@ -0,0 +1,23 @@ +using System.Runtime.CompilerServices; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.Implementations; + +public class SchemaRegistryDefinitionEnumerator(IEnumerable providers) : ISchemaRegistryDefinitionEnumerator +{ + public async Task> EnumerateAsync(CancellationToken cancellationToken) + { + return await GetSchemaRegistryDefinitionsInternalAsync(cancellationToken).ToListAsync(cancellationToken); + } + + private async IAsyncEnumerable GetSchemaRegistryDefinitionsInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + foreach (var provider in providers) + { + var schemaRegistryDefinitions = await provider.GetSchemaRegistryDefinitionsAsync(cancellationToken).ToList(); + + foreach (var schemaRegistryDefinition in schemaRegistryDefinitions) + yield return schemaRegistryDefinition; + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Implementations/WorkerManager.cs b/src/modules/Elsa.Kafka/Implementations/WorkerManager.cs index a4d7b0d70..b3383d123 100644 --- a/src/modules/Elsa.Kafka/Implementations/WorkerManager.cs +++ b/src/modules/Elsa.Kafka/Implementations/WorkerManager.cs @@ -41,7 +41,7 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) : if (existingWorker == null) { // Add a new worker. - var worker = CreateWorker(scope.ServiceProvider, consumerDefinition); + var worker = await CreateWorkerAsync(scope.ServiceProvider, consumerDefinition, cancellationToken); workers.Add(consumerDefinition.Id, worker); } else @@ -54,7 +54,7 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) : if (existingHash != hash) { existingWorker.Stop(); - var worker = CreateWorker(scope.ServiceProvider, consumerDefinition); + var worker = await CreateWorkerAsync(scope.ServiceProvider, consumerDefinition, cancellationToken); workers[consumerDefinition.Id] = worker; } } @@ -162,7 +162,7 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) : return Workers.TryGetValue(consumerDefinitionId, out var worker) ? worker : null; } - private IWorker CreateWorker(IServiceProvider serviceProvider, ConsumerDefinition consumerDefinition) + private async Task CreateWorkerAsync(IServiceProvider serviceProvider, ConsumerDefinition consumerDefinition, CancellationToken cancellationToken) { var factoryType = consumerDefinition.FactoryType; var consumerFactory = ActivatorUtilities.GetServiceOrCreateInstance(serviceProvider, factoryType) as IConsumerFactory; @@ -170,7 +170,8 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) : if (consumerFactory == null) throw new InvalidOperationException($"Worker factory of type '{factoryType}' not found."); - var createConsumerContext = new CreateConsumerContext(consumerDefinition); + var schemaRegistryDefinition = await GetSchemaRegistryDefinitionAsync(serviceProvider, consumerDefinition.SchemaRegistryId, cancellationToken); + var createConsumerContext = new CreateConsumerContext(consumerDefinition, schemaRegistryDefinition); var consumerProxy = consumerFactory.CreateConsumer(createConsumerContext); var wrappedConsumer = consumerProxy.Consumer; var wrappedConsumerType = wrappedConsumer.GetType(); @@ -186,4 +187,14 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) : { return hasher.Hash(consumerDefinition); } + + private async Task GetSchemaRegistryDefinitionAsync(IServiceProvider serviceProvider, string? id, CancellationToken cancellationToken = default) + { + if (id == null) + return null; + + var schemaRegistryDefinitionEnumerator = serviceProvider.GetRequiredService(); + var schemaRegistryDefinitions = await schemaRegistryDefinitionEnumerator.EnumerateAsync(cancellationToken).ToList(); + return schemaRegistryDefinitions.FirstOrDefault(x => x.Id == id); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Models/CreateConsumerContext.cs b/src/modules/Elsa.Kafka/Models/CreateConsumerContext.cs index 1d2bcc65d..c50b69b3c 100644 --- a/src/modules/Elsa.Kafka/Models/CreateConsumerContext.cs +++ b/src/modules/Elsa.Kafka/Models/CreateConsumerContext.cs @@ -1,3 +1,3 @@ namespace Elsa.Kafka; -public record CreateConsumerContext(ConsumerDefinition ConsumerDefinition); \ No newline at end of file +public record CreateConsumerContext(ConsumerDefinition ConsumerDefinition, SchemaRegistryDefinition? SchemaRegistryDefinition); \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Models/CreateProducerContext.cs b/src/modules/Elsa.Kafka/Models/CreateProducerContext.cs index ca1edb364..4b76ee286 100644 --- a/src/modules/Elsa.Kafka/Models/CreateProducerContext.cs +++ b/src/modules/Elsa.Kafka/Models/CreateProducerContext.cs @@ -1,3 +1,3 @@ namespace Elsa.Kafka; -public record CreateProducerContext(ProducerDefinition ProducerDefinition); \ No newline at end of file +public record CreateProducerContext(ProducerDefinition ProducerDefinition, SchemaRegistryDefinition? SchemaRegistryDefinition); \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Options/KafkaOptions.cs b/src/modules/Elsa.Kafka/Options/KafkaOptions.cs index 32d99aec1..895b590e3 100644 --- a/src/modules/Elsa.Kafka/Options/KafkaOptions.cs +++ b/src/modules/Elsa.Kafka/Options/KafkaOptions.cs @@ -3,6 +3,7 @@ namespace Elsa.Kafka; public class KafkaOptions { public ICollection Topics { get; set; } = []; + public ICollection SchemaRegistries { get; set; } = []; public ICollection Consumers { get; set; } = []; public ICollection Producers { get; set; } = []; public string WorkflowInstanceIdHeaderKey { get; set; } = "x-workflow-instance-id"; diff --git a/src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs b/src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs index 668146b5b..30677c184 100644 --- a/src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs +++ b/src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs @@ -2,7 +2,7 @@ using Microsoft.Extensions.Options; namespace Elsa.Kafka.Providers; -public class OptionsDefinitionProvider(IOptions options) : IConsumerDefinitionProvider, IProducerDefinitionProvider, ITopicDefinitionProvider +public class OptionsDefinitionProvider(IOptions options) : IConsumerDefinitionProvider, IProducerDefinitionProvider, ITopicDefinitionProvider, ISchemaRegistryDefinitionProvider { public Task> GetConsumerDefinitionsAsync(CancellationToken cancellationToken = default) { @@ -18,4 +18,9 @@ public class OptionsDefinitionProvider(IOptions options) : IConsum { return Task.FromResult>(options.Value.Topics); } + + public Task> GetSchemaRegistryDefinitionsAsync(CancellationToken cancellationToken = default) + { + return Task.FromResult>(options.Value.SchemaRegistries); + } } \ No newline at end of file