From 6022df165cf63b88256652ba328bc0febeee89bc Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 29 Nov 2024 19:49:31 +0100 Subject: [PATCH] Kafka: Update ProduceMessage activity with support for specifying a Key (#6166) * Add Key to Kafka ProduceMessage activity Deleted unnecessary Consumer and Producer workflow classes and the OrderReceived message class to clean up code. Refactored Kafka producer interface and implementation to include message keys for improved message handling. Updated configuration to enable Kafka and removed unused service registrations. * Add Kafka factory classes and type alias registry Introduce GenericConsumerFactory and GenericProducerFactory for handling Kafka consumer and producer creation. Implement a TypeAliasRegistry to manage type aliases, enabling cleaner configuration through aliases. Update the OrderReceived message class and ensure better integration with the server web program via these new components. * Handle empty topics and predicates in Kafka worker. Ensure the Kafka consumer unsubscribes when no topics are available to subscribe to. Additionally, add a check to handle empty string values for predicates, allowing workflow triggers to proceed in this scenario. * Disable Kafka usage in Elsa Server Web configuration Kafka has been disabled in the current configuration by setting the useKafka constant to false. This change might be intended to switch to a different messaging system or to simplify the current setup by removing unnecessary services. Ensure that any dependencies on Kafka are handled elsewhere in the application. --- .../Elsa.Server.Web/Messages/OrderReceived.cs | 1 - src/apps/Elsa.Server.Web/Program.cs | 13 +++++--- .../ConsumerConfigWorkflowContextProvider.cs | 22 ------------- .../Workflows/ConsumerWorkflow.cs | 32 ------------------- .../Workflows/ProducerWorkflow.cs | 22 ------------- src/apps/Elsa.Server.Web/appsettings.json | 12 ++----- .../Serialization/TypeAliasRegistry.cs | 10 ++++++ .../Serialization/TypeTypeConverter.cs | 8 +++++ .../Elsa.Kafka/Activities/ProduceMessage.cs | 31 +++++++++++------- src/modules/Elsa.Kafka/Contracts/IProducer.cs | 2 +- .../Factories/GenericConsumerFactory.cs | 16 ++++++++++ .../Factories/GenericProducerFactory.cs | 16 ++++++++++ .../Elsa.Kafka/Handlers/TriggerWorkflows.cs | 3 ++ .../Implementations/ProducerProxy.cs | 13 ++++---- .../Elsa.Kafka/Implementations/Worker.cs | 5 ++- .../Implementations/WorkerManager.cs | 5 +-- 16 files changed, 97 insertions(+), 114 deletions(-) delete mode 100644 src/apps/Elsa.Server.Web/WorkflowContextProviders/ConsumerConfigWorkflowContextProvider.cs delete mode 100644 src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs delete mode 100644 src/apps/Elsa.Server.Web/Workflows/ProducerWorkflow.cs create mode 100644 src/modules/Elsa.Common/Serialization/TypeAliasRegistry.cs create mode 100644 src/modules/Elsa.Kafka/Factories/GenericConsumerFactory.cs create mode 100644 src/modules/Elsa.Kafka/Factories/GenericProducerFactory.cs diff --git a/src/apps/Elsa.Server.Web/Messages/OrderReceived.cs b/src/apps/Elsa.Server.Web/Messages/OrderReceived.cs index 350c988b8..411081bd1 100644 --- a/src/apps/Elsa.Server.Web/Messages/OrderReceived.cs +++ b/src/apps/Elsa.Server.Web/Messages/OrderReceived.cs @@ -3,5 +3,4 @@ namespace Elsa.Server.Web.Messages; public class OrderReceived { public string OrderId { get; set; } = default!; - public decimal OrderTotal { get; set; } = default!; } \ No newline at end of file diff --git a/src/apps/Elsa.Server.Web/Program.cs b/src/apps/Elsa.Server.Web/Program.cs index 6c80b6fce..129a232d2 100644 --- a/src/apps/Elsa.Server.Web/Program.cs +++ b/src/apps/Elsa.Server.Web/Program.cs @@ -4,6 +4,7 @@ using Elsa.Alterations.Extensions; using Elsa.Alterations.MassTransit.Extensions; using Elsa.Common.DistributedHosting.DistributedLocks; using Elsa.Common.RecurringTasks; +using Elsa.Common.Serialization; using Elsa.Dapper.Extensions; using Elsa.Dapper.Services; using Elsa.DropIns.Extensions; @@ -16,6 +17,7 @@ using Elsa.Extensions; using Elsa.Features.Services; using Elsa.Identity.Multitenancy; using Elsa.Kafka; +using Elsa.Kafka.Factories; using Elsa.MassTransit.Extensions; using Elsa.MongoDb.Extensions; using Elsa.MongoDb.Modules.Alterations; @@ -30,11 +32,11 @@ using Elsa.Server.Web; using Elsa.Server.Web.Extensions; using Elsa.Server.Web.Filters; using Elsa.Server.Web.Messages; -using Elsa.Server.Web.WorkflowContextProviders; using Elsa.Tenants.AspNetCore; using Elsa.Tenants.Extensions; using Elsa.Workflows.Api; using Elsa.Workflows.LogPersistence; +using Elsa.Workflows.Management; using Elsa.Workflows.Management.Compression; using Elsa.Workflows.Management.Stores; using Elsa.Workflows.Runtime.Distributed.Extensions; @@ -93,6 +95,10 @@ var redisConnectionString = configuration.GetConnectionString("Redis")!; var distributedLockProviderName = configuration.GetSection("Runtime:DistributedLocking")["Provider"]; var appRole = Enum.Parse(configuration["AppRole"] ?? "Default"); +// Optionally create type aliases for easier configuration. +TypeAliasRegistry.RegisterAlias("OrderReceivedProducerFactory", typeof(GenericProducerFactory)); +TypeAliasRegistry.RegisterAlias("OrderReceivedConsumerFactory", typeof(GenericConsumerFactory)); + // Add Elsa services. services .AddElsa(elsa => @@ -198,6 +204,7 @@ services management.SetDefaultLogPersistenceMode(LogPersistenceMode.Inherit); management.UseReadOnlyMode(useReadOnlyMode); + management.AddVariableTypeAndAlias("Application"); }) .UseProtoActor(proto => { @@ -426,8 +433,6 @@ services // etc. }); } - - massTransit.AddMessageType(); }); } @@ -454,8 +459,6 @@ services { kafka.ConfigureOptions(options => configuration.GetSection("Kafka").Bind(options)); }); - - services.AddWorkflowContextProvider(); } if (useAgents) diff --git a/src/apps/Elsa.Server.Web/WorkflowContextProviders/ConsumerConfigWorkflowContextProvider.cs b/src/apps/Elsa.Server.Web/WorkflowContextProviders/ConsumerConfigWorkflowContextProvider.cs deleted file mode 100644 index 5acf2f7c6..000000000 --- a/src/apps/Elsa.Server.Web/WorkflowContextProviders/ConsumerConfigWorkflowContextProvider.cs +++ /dev/null @@ -1,22 +0,0 @@ -using Elsa.Common.Multitenancy; -using Elsa.Kafka; -using Elsa.WorkflowContexts.Abstractions; -using Elsa.Workflows; - -namespace Elsa.Server.Web.WorkflowContextProviders; - -public class ConsumerDefinitionWorkflowContextProvider(IConsumerDefinitionEnumerator consumerDefinitionEnumerator, ITenantAccessor tenantAccessor) : WorkflowContextProvider -{ - protected override async ValueTask LoadAsync(WorkflowExecutionContext workflowExecutionContext) - { - var tenant = tenantAccessor.Tenant; - var tenantId = tenant?.Id; - var definitionId = workflowExecutionContext.Workflow.Identity.DefinitionId; - - // Load specific setting here. - - // For now, just return the first consumer definition. - var consumerDefinitions = await consumerDefinitionEnumerator.EnumerateAsync(workflowExecutionContext.CancellationToken); - return consumerDefinitions.FirstOrDefault(); - } -} \ No newline at end of file diff --git a/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs b/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs deleted file mode 100644 index cf24baeb3..000000000 --- a/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs +++ /dev/null @@ -1,32 +0,0 @@ -using System.Dynamic; -using System.Text.Json; -using Elsa.JavaScript.Models; -using Elsa.Kafka.Activities; -using Elsa.Workflows; -using Elsa.Workflows.Activities; - -namespace Elsa.Server.Web.Workflows; - -public class ConsumerWorkflow : WorkflowBase -{ - protected override void Build(IWorkflowBuilder builder) - { - var message = builder.WithVariable(); - builder.Name = "Consumer Workflow"; - builder.Root = new Sequence - { - Activities = - { - new MessageReceived - { - ConsumerDefinitionId = new("consumer-1"), - Topics = new(["topic-1"]), - Predicate = new(JavaScriptExpression.Create("getMessage().OrderId == '1'")), - Result = new(message), - CanStartWorkflow = true - }, - new WriteLine(c => JsonSerializer.Serialize(message.Get(c))) - } - }; - } -} \ No newline at end of file diff --git a/src/apps/Elsa.Server.Web/Workflows/ProducerWorkflow.cs b/src/apps/Elsa.Server.Web/Workflows/ProducerWorkflow.cs deleted file mode 100644 index 42624c2e1..000000000 --- a/src/apps/Elsa.Server.Web/Workflows/ProducerWorkflow.cs +++ /dev/null @@ -1,22 +0,0 @@ -using Elsa.Kafka.Activities; -using Elsa.Workflows; - -namespace Elsa.Server.Web.Workflows; - -public class ProducerWorkflow : WorkflowBase -{ - protected override void Build(IWorkflowBuilder builder) - { - builder.Name = "Producer Workflow"; - builder.Root = new ProduceMessage - { - Topic = new("topic-2"), - ProducerDefinitionId = new("producer-1"), - Content = new(() => new - { - OrderId = "1", - CustomerId = "1" - }) - }; - } -} \ No newline at end of file diff --git a/src/apps/Elsa.Server.Web/appsettings.json b/src/apps/Elsa.Server.Web/appsettings.json index 14daa0ed0..79f00e793 100644 --- a/src/apps/Elsa.Server.Web/appsettings.json +++ b/src/apps/Elsa.Server.Web/appsettings.json @@ -206,21 +206,13 @@ { "Id": "topic-2", "Name": "topic-2" - }, - { - "Id": "topic-3", - "Name": "topic-3" - }, - { - "Id": "topic-4", - "Name": "topic-4" } ], "Producers": [ { "Id": "producer-1", "Name": "Producer 1", - "FactoryType": "Elsa.Kafka.Factories.ExpandoObjectProducerFactory, Elsa.Kafka", + "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" } @@ -230,7 +222,7 @@ { "Id": "consumer-1", "Name": "Consumer 1", - "FactoryType": "Elsa.Kafka.Factories.ExpandoObjectConsumerFactory, Elsa.Kafka", + "FactoryType": "Elsa.Kafka.Factories.GenericConsumerFactory`2[[System.String, System.Private.CoreLib], [Elsa.Server.Web.Messages.OrderReceived, Elsa.Server.Web]], Elsa.Kafka", "Config": { "BootstrapServers": "localhost:9092", "GroupId": "group-1", diff --git a/src/modules/Elsa.Common/Serialization/TypeAliasRegistry.cs b/src/modules/Elsa.Common/Serialization/TypeAliasRegistry.cs new file mode 100644 index 000000000..7c610ef41 --- /dev/null +++ b/src/modules/Elsa.Common/Serialization/TypeAliasRegistry.cs @@ -0,0 +1,10 @@ +namespace Elsa.Common.Serialization; + +public static class TypeAliasRegistry +{ + public static Dictionary TypeAliases { get; } = new(); + + public static void RegisterAlias(string alias, Type type) => TypeAliases[alias] = type; + + public static Type? GetType(string alias) => TypeAliases.GetValueOrDefault(alias); +} \ No newline at end of file diff --git a/src/modules/Elsa.Common/Serialization/TypeTypeConverter.cs b/src/modules/Elsa.Common/Serialization/TypeTypeConverter.cs index 29ea25f03..4d89b6c5e 100644 --- a/src/modules/Elsa.Common/Serialization/TypeTypeConverter.cs +++ b/src/modules/Elsa.Common/Serialization/TypeTypeConverter.cs @@ -15,7 +15,11 @@ public class TypeTypeConverter : TypeConverter public override object? ConvertFrom(ITypeDescriptorContext? context, CultureInfo? culture, object value) { if (value is string stringValue) + { + if (TypeAliasRegistry.GetType(stringValue) is { } type) + return type; return Type.GetType(stringValue); + } return base.ConvertFrom(context, culture, value); } @@ -27,7 +31,11 @@ public class TypeTypeConverter : TypeConverter public override object? ConvertTo(ITypeDescriptorContext? context, CultureInfo? culture, object? value, Type destinationType) { if (destinationType == typeof(string) && value is Type type) + { + if (TypeAliasRegistry.TypeAliases.FirstOrDefault(x => x.Value == type).Key is { } alias) + return alias; return type.AssemblyQualifiedName; + } return base.ConvertTo(context, culture, value, destinationType); } } \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs b/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs index a5e2b1ed6..381c423b1 100644 --- a/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs +++ b/src/modules/Elsa.Kafka/Activities/ProduceMessage.cs @@ -7,6 +7,7 @@ using Elsa.Workflows.Attributes; using Elsa.Workflows.Models; using Elsa.Workflows.Runtime; using Elsa.Workflows.UIHints; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; namespace Elsa.Kafka.Activities; @@ -48,25 +49,33 @@ public class ProduceMessage : CodeActivity public Input CorrelationId { get; set; } = default!; /// - /// The content of the message to send. + /// The content of the message to produce. /// [Input(Description = "The content of the message to produce.")] public Input Content { get; set; } = default!; + /// + /// The key of the message to send. + /// + [Input(Description = "The key of the message to produce.")] + public Input Key { get; set; } = default!; + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { + var cancellationToken = context.CancellationToken; var topic = Topic.Get(context); var producerDefinitionId = ProducerDefinitionId.Get(context); var producerDefinitionEnumerator = context.GetRequiredService(); var producerDefinition = await producerDefinitionEnumerator.GetByIdAsync(producerDefinitionId); var content = Content.Get(context); + var key = Key.Get(context); - context.DeferTask(async () => - { - using var producer = CreateProducer(context, producerDefinition); - var headers = CreateHeaders(context); - await producer.ProduceAsync(topic, content, headers); - }); + if (key is string keyString && string.IsNullOrWhiteSpace(keyString)) + key = null; + + using var producer = CreateProducer(context, producerDefinition); + var headers = CreateHeaders(context); + await producer.ProduceAsync(topic, key, content, headers, cancellationToken); } private Headers CreateHeaders(ActivityExecutionContext context) @@ -84,14 +93,14 @@ public class ProduceMessage : CodeActivity return headers; } - + private IProducer CreateProducer(ActivityExecutionContext context, ProducerDefinition producerDefinition) { - var factory = context.GetRequiredService(producerDefinition.FactoryType) as IProducerFactory; - + 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); return factory.CreateProducer(createProducerContext); } diff --git a/src/modules/Elsa.Kafka/Contracts/IProducer.cs b/src/modules/Elsa.Kafka/Contracts/IProducer.cs index 740195854..9ab4e4b6b 100644 --- a/src/modules/Elsa.Kafka/Contracts/IProducer.cs +++ b/src/modules/Elsa.Kafka/Contracts/IProducer.cs @@ -4,5 +4,5 @@ namespace Elsa.Kafka; public interface IProducer : IDisposable { - Task ProduceAsync(string topic, object value, Headers? headers = null, CancellationToken cancellationToken = default); + Task ProduceAsync(string topic, object? key, object value, Headers? headers = null, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Factories/GenericConsumerFactory.cs b/src/modules/Elsa.Kafka/Factories/GenericConsumerFactory.cs new file mode 100644 index 000000000..e1ea67e83 --- /dev/null +++ b/src/modules/Elsa.Kafka/Factories/GenericConsumerFactory.cs @@ -0,0 +1,16 @@ +using Confluent.Kafka; +using Elsa.Kafka.Implementations; +using Elsa.Kafka.Serializers; + +namespace Elsa.Kafka.Factories; + +public class GenericConsumerFactory : IConsumerFactory +{ + public IConsumer CreateConsumer(CreateConsumerContext context) + { + var consumer = new ConsumerBuilder(context.ConsumerDefinition.Config) + .SetValueDeserializer(new JsonDeserializer()) + .Build(); + return new ConsumerProxy(consumer); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Factories/GenericProducerFactory.cs b/src/modules/Elsa.Kafka/Factories/GenericProducerFactory.cs new file mode 100644 index 000000000..2f97a951e --- /dev/null +++ b/src/modules/Elsa.Kafka/Factories/GenericProducerFactory.cs @@ -0,0 +1,16 @@ +using Confluent.Kafka; +using Elsa.Kafka.Implementations; +using Elsa.Kafka.Serializers; + +namespace Elsa.Kafka.Factories; + +public class GenericProducerFactory : IProducerFactory +{ + public IProducer CreateProducer(CreateProducerContext workerContext) + { + var producer = new ProducerBuilder(workerContext.ProducerDefinition.Config) + .SetValueSerializer(new JsonSerializer()) + .Build(); + return new ProducerProxy(producer); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs b/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs index c85659947..af610d7da 100644 --- a/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs +++ b/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs @@ -183,6 +183,9 @@ public class TriggerWorkflows( if (predicate == null) return true; + + if(string.IsNullOrWhiteSpace(predicate.Value as string)) + return true; var expressionExecutionContext = await GetExpressionExecutionContextAsync(transportMessage, binding, cancellationToken); diff --git a/src/modules/Elsa.Kafka/Implementations/ProducerProxy.cs b/src/modules/Elsa.Kafka/Implementations/ProducerProxy.cs index 9741298db..321b06601 100644 --- a/src/modules/Elsa.Kafka/Implementations/ProducerProxy.cs +++ b/src/modules/Elsa.Kafka/Implementations/ProducerProxy.cs @@ -7,7 +7,7 @@ public class ProducerProxy(object producer) : IProducer { private object Producer { get; } = producer; - public async Task ProduceAsync(string topic, object value, Headers? headers = null, CancellationToken cancellationToken = default) + public async Task ProduceAsync(string topic, object? key, object value, Headers? headers = null, CancellationToken cancellationToken = default) { var producerType = Producer.GetType(); var keyType = producerType.GetGenericArguments()[0]; @@ -16,14 +16,13 @@ public class ProducerProxy(object producer) : IProducer var produceAsyncMethod = producerType.GetMethod("ProduceAsync", [typeof(string), messageType, typeof(CancellationToken)])!; var messageInstance = Activator.CreateInstance(messageType); var convertedValue = value.ConvertTo(valueType); - + messageType.GetProperty("Value")!.SetValue(messageInstance, convertedValue); - - if (headers != null) - messageType.GetProperty("Headers")!.SetValue(messageInstance, headers); - + if (key != null) messageType.GetProperty("Key")!.SetValue(messageInstance, key); + if (headers != null) messageType.GetProperty("Headers")!.SetValue(messageInstance, headers); + await (Task)produceAsyncMethod.Invoke(Producer, [topic, messageInstance, cancellationToken])!; - + var flushMethod = producerType.GetMethod("Flush", [typeof(CancellationToken)])!; flushMethod.Invoke(Producer, [cancellationToken]); } diff --git a/src/modules/Elsa.Kafka/Implementations/Worker.cs b/src/modules/Elsa.Kafka/Implementations/Worker.cs index 218cdf23a..fa3a7825b 100644 --- a/src/modules/Elsa.Kafka/Implementations/Worker.cs +++ b/src/modules/Elsa.Kafka/Implementations/Worker.cs @@ -90,7 +90,10 @@ public class Worker(WorkerContext workerContext, IConsumer