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