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