diff --git a/Directory.Packages.props b/Directory.Packages.props index f994e73c3..7c500f677 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -15,6 +15,7 @@ + @@ -84,7 +85,7 @@ - + diff --git a/Elsa.sln b/Elsa.sln index 19e4b2aa6..619e260be 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -91,6 +91,7 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "docker", "docker", "{986E54 docker\ElsaStudio.Dockerfile = docker\ElsaStudio.Dockerfile docker\init-db.sh = docker\init-db.sh docker\otel-collector-config.yaml = docker\otel-collector-config.yaml + docker\docker-compose-kafka.yml = docker\docker-compose-kafka.yml EndProjectSection EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Elasticsearch", "src\modules\Elsa.Elasticsearch\Elsa.Elasticsearch.csproj", "{3246883E-2FA7-4B4A-BDC5-99039A2869BC}" @@ -384,6 +385,8 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Tenants.AspNetCore", " EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Agents.Persistence.EntityFrameworkCore.PostgreSql", "src\modules\Elsa.Agents.Persistence.EntityFrameworkCore.PostgreSql\Elsa.Agents.Persistence.EntityFrameworkCore.PostgreSql.csproj", "{2B939AC9-03A4-479E-AA0D-CB58F4A7F480}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Kafka", "src\modules\Elsa.Kafka\Elsa.Kafka.csproj", "{BF934627-F531-44FB-BEC2-ECA801FF31E7}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -816,6 +819,10 @@ Global {2B939AC9-03A4-479E-AA0D-CB58F4A7F480}.Debug|Any CPU.Build.0 = Debug|Any CPU {2B939AC9-03A4-479E-AA0D-CB58F4A7F480}.Release|Any CPU.ActiveCfg = Release|Any CPU {2B939AC9-03A4-479E-AA0D-CB58F4A7F480}.Release|Any CPU.Build.0 = Release|Any CPU + {BF934627-F531-44FB-BEC2-ECA801FF31E7}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {BF934627-F531-44FB-BEC2-ECA801FF31E7}.Debug|Any CPU.Build.0 = Debug|Any CPU + {BF934627-F531-44FB-BEC2-ECA801FF31E7}.Release|Any CPU.ActiveCfg = Release|Any CPU + {BF934627-F531-44FB-BEC2-ECA801FF31E7}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -960,6 +967,7 @@ Global {D5720DBC-8C2B-42D5-9D9F-2FF6EAD4001C} = {2F3E1026-5054-4E1F-899B-F1A7F70F9912} {2B939AC9-03A4-479E-AA0D-CB58F4A7F480} = {50470834-4CD8-479A-8B58-0A1869BA5D37} {2CDF3E1C-267D-4198-B1C7-7E1F548FC120} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} + {BF934627-F531-44FB-BEC2-ECA801FF31E7} = {DD089B8B-DA73-492A-9010-F772D1C178DA} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {D4B5CEAA-7D70-4FCB-A68E-B03FBE5E0E5E} diff --git a/docker/docker-compose-kafka.yml b/docker/docker-compose-kafka.yml new file mode 100644 index 000000000..106510d77 --- /dev/null +++ b/docker/docker-compose-kafka.yml @@ -0,0 +1,167 @@ +--- +version: '2' +services: + + broker: + image: confluentinc/cp-kafka:7.7.1 + hostname: broker + container_name: broker + ports: + - "9092:9092" + - "9101:9101" + environment: + KAFKA_NODE_ID: 1 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT' + KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092' + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_JMX_PORT: 9101 + KAFKA_JMX_HOSTNAME: localhost + KAFKA_PROCESS_ROLES: 'broker,controller' + KAFKA_CONTROLLER_QUORUM_VOTERS: '1@broker:29093' + KAFKA_LISTENERS: 'PLAINTEXT://broker:29092,CONTROLLER://broker:29093,PLAINTEXT_HOST://0.0.0.0:9092' + KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT' + KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' + KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' + # Replace CLUSTER_ID with a unique base64 UUID using "bin/kafka-storage.sh random-uuid" + # See https://docs.confluent.io/kafka/operations-tools/kafka-tools.html#kafka-storage-sh + CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk' + + schema-registry: + image: confluentinc/cp-schema-registry:7.7.1 + hostname: schema-registry + container_name: schema-registry + depends_on: + - broker + ports: + - "8081:8081" + environment: + SCHEMA_REGISTRY_HOST_NAME: schema-registry + SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: 'broker:29092' + SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 + + connect: + image: cnfldemos/cp-server-connect-datagen:0.6.4-7.6.0 + hostname: connect + container_name: connect + depends_on: + - broker + - schema-registry + ports: + - "8083:8083" + environment: + CONNECT_BOOTSTRAP_SERVERS: 'broker:29092' + CONNECT_REST_ADVERTISED_HOST_NAME: connect + CONNECT_GROUP_ID: compose-connect-group + CONNECT_CONFIG_STORAGE_TOPIC: docker-connect-configs + CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1 + CONNECT_OFFSET_FLUSH_INTERVAL_MS: 10000 + CONNECT_OFFSET_STORAGE_TOPIC: docker-connect-offsets + CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1 + CONNECT_STATUS_STORAGE_TOPIC: docker-connect-status + CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1 + CONNECT_KEY_CONVERTER: org.apache.kafka.connect.storage.StringConverter + CONNECT_VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter + CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081 + # CLASSPATH required due to CC-2422 + CLASSPATH: /usr/share/java/monitoring-interceptors/monitoring-interceptors-7.7.1.jar + CONNECT_PRODUCER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringProducerInterceptor" + CONNECT_CONSUMER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringConsumerInterceptor" + CONNECT_PLUGIN_PATH: "/usr/share/java,/usr/share/confluent-hub-components" + CONNECT_LOG4J_LOGGERS: org.apache.zookeeper=ERROR,org.I0Itec.zkclient=ERROR,org.reflections=ERROR + + control-center: + image: confluentinc/cp-enterprise-control-center:7.7.1 + hostname: control-center + container_name: control-center + depends_on: + - broker + - schema-registry + - connect + - ksqldb-server + ports: + - "9021:9021" + environment: + CONTROL_CENTER_BOOTSTRAP_SERVERS: 'broker:29092' + CONTROL_CENTER_CONNECT_CONNECT-DEFAULT_CLUSTER: 'connect:8083' + CONTROL_CENTER_CONNECT_HEALTHCHECK_ENDPOINT: '/connectors' + CONTROL_CENTER_KSQL_KSQLDB1_URL: "http://ksqldb-server:8088" + CONTROL_CENTER_KSQL_KSQLDB1_ADVERTISED_URL: "http://localhost:8088" + CONTROL_CENTER_SCHEMA_REGISTRY_URL: "http://schema-registry:8081" + CONTROL_CENTER_REPLICATION_FACTOR: 1 + CONTROL_CENTER_INTERNAL_TOPICS_PARTITIONS: 1 + CONTROL_CENTER_MONITORING_INTERCEPTOR_TOPIC_PARTITIONS: 1 + CONFLUENT_METRICS_TOPIC_REPLICATION: 1 + PORT: 9021 + + ksqldb-server: + image: confluentinc/cp-ksqldb-server:7.7.1 + hostname: ksqldb-server + container_name: ksqldb-server + depends_on: + - broker + - connect + ports: + - "8088:8088" + environment: + KSQL_CONFIG_DIR: "/etc/ksql" + KSQL_BOOTSTRAP_SERVERS: "broker:29092" + KSQL_HOST_NAME: ksqldb-server + KSQL_LISTENERS: "http://0.0.0.0:8088" + KSQL_CACHE_MAX_BYTES_BUFFERING: 0 + KSQL_KSQL_SCHEMA_REGISTRY_URL: "http://schema-registry:8081" + KSQL_PRODUCER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringProducerInterceptor" + KSQL_CONSUMER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringConsumerInterceptor" + KSQL_KSQL_CONNECT_URL: "http://connect:8083" + KSQL_KSQL_LOGGING_PROCESSING_TOPIC_REPLICATION_FACTOR: 1 + KSQL_KSQL_LOGGING_PROCESSING_TOPIC_AUTO_CREATE: 'true' + KSQL_KSQL_LOGGING_PROCESSING_STREAM_AUTO_CREATE: 'true' + + ksqldb-cli: + image: confluentinc/cp-ksqldb-cli:7.7.1 + container_name: ksqldb-cli + depends_on: + - broker + - connect + - ksqldb-server + entrypoint: /bin/sh + tty: true + + ksql-datagen: + image: confluentinc/ksqldb-examples:7.7.1 + hostname: ksql-datagen + container_name: ksql-datagen + depends_on: + - ksqldb-server + - broker + - schema-registry + - connect + command: "bash -c 'echo Waiting for Kafka to be ready... && \ + cub kafka-ready -b broker:29092 1 40 && \ + echo Waiting for Confluent Schema Registry to be ready... && \ + cub sr-ready schema-registry 8081 40 && \ + echo Waiting a few seconds for topic creation to finish... && \ + sleep 11 && \ + tail -f /dev/null'" + environment: + KSQL_CONFIG_DIR: "/etc/ksql" + STREAMS_BOOTSTRAP_SERVERS: broker:29092 + STREAMS_SCHEMA_REGISTRY_HOST: schema-registry + STREAMS_SCHEMA_REGISTRY_PORT: 8081 + + rest-proxy: + image: confluentinc/cp-kafka-rest:7.7.1 + depends_on: + - broker + - schema-registry + ports: + - 8082:8082 + hostname: rest-proxy + container_name: rest-proxy + environment: + KAFKA_REST_HOST_NAME: rest-proxy + KAFKA_REST_BOOTSTRAP_SERVERS: 'broker:29092' + KAFKA_REST_LISTENERS: "http://0.0.0.0:8082" + KAFKA_REST_SCHEMA_REGISTRY_URL: 'http://schema-registry:8081' diff --git a/docker/docker-compose.yml b/docker/docker-compose.yml index 51cb37e4e..40ee95d18 100644 --- a/docker/docker-compose.yml +++ b/docker/docker-compose.yml @@ -11,6 +11,19 @@ - ./init-db.sh:/docker-entrypoint-initdb.d/init-db.sh # This will initialize the 'tracelens' database ports: - "5432:5432" + + mysql: + image: mysql:9.1.0 + container_name: mysql + environment: + MYSQL_ROOT_PASSWORD: password + MYSQL_DATABASE: elsa + MYSQL_USER: admin + MYSQL_PASSWORD: password + ports: + - "3306:3306" + volumes: + - mysql_data:/var/lib/mysql mongodb: image: mongo:latest @@ -91,5 +104,6 @@ volumes: postgres-data: + mysql_data: cockroachdb-data: mongodb_data: diff --git a/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj b/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj index accb12903..0b5e96890 100644 --- a/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj +++ b/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj @@ -3,13 +3,17 @@ + + + + diff --git a/src/apps/Elsa.Server.Web/Program.cs b/src/apps/Elsa.Server.Web/Program.cs index e3d33877c..1f3d9e85d 100644 --- a/src/apps/Elsa.Server.Web/Program.cs +++ b/src/apps/Elsa.Server.Web/Program.cs @@ -16,6 +16,7 @@ using Elsa.EntityFrameworkCore.Modules.Runtime; using Elsa.Extensions; using Elsa.Features.Services; using Elsa.Identity.Multitenancy; +using Elsa.Kafka; using Elsa.MassTransit.Extensions; using Elsa.MongoDb.Extensions; using Elsa.MongoDb.Modules.Alterations; @@ -30,6 +31,7 @@ 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; @@ -63,6 +65,7 @@ const bool runEFCoreMigrations = true; const bool useMemoryStores = false; const bool useCaching = true; const bool useAzureServiceBus = false; +const bool useKafka = true; const bool useReadOnlyMode = false; const bool useSignalR = false; // Disabled until Elsa Studio sends authenticated requests. const WorkflowRuntime workflowRuntime = WorkflowRuntime.ProtoActor; @@ -81,6 +84,7 @@ var identityTokenSection = identitySection.GetSection("Tokens"); var sqliteConnectionString = configuration.GetConnectionString("Sqlite")!; var sqlServerConnectionString = configuration.GetConnectionString("SqlServer")!; var postgresConnectionString = configuration.GetConnectionString("PostgreSql")!; +var mySqlConnectionString = configuration.GetConnectionString("MySql")!; var cockroachDbConnectionString = configuration.GetConnectionString("CockroachDb")!; var mongoDbConnectionString = configuration.GetConnectionString("MongoDb")!; var azureServiceBusConnectionString = configuration.GetConnectionString("AzureServiceBus")!; @@ -136,6 +140,8 @@ services ef.UseSqlServer(sqlServerConnectionString!); else if (sqlDatabaseProvider == SqlDatabaseProvider.PostgreSql) ef.UsePostgreSql(postgresConnectionString!); + else if(sqlDatabaseProvider == SqlDatabaseProvider.MySql) + ef.UseMySql(mySqlConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.CockroachDb) ef.UsePostgreSql(cockroachDbConnectionString!); else @@ -168,6 +174,8 @@ services ef.UseSqlServer(sqlServerConnectionString!); else if (sqlDatabaseProvider == SqlDatabaseProvider.PostgreSql) ef.UsePostgreSql(postgresConnectionString!); + else if(sqlDatabaseProvider == SqlDatabaseProvider.MySql) + ef.UseMySql(mySqlConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.CockroachDb) ef.UsePostgreSql(cockroachDbConnectionString!); else @@ -236,6 +244,8 @@ services } else if (sqlDatabaseProvider == SqlDatabaseProvider.PostgreSql) ef.UsePostgreSql(postgresConnectionString!); + else if(sqlDatabaseProvider == SqlDatabaseProvider.MySql) + ef.UseMySql(mySqlConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.CockroachDb) ef.UsePostgreSql(cockroachDbConnectionString!); else @@ -359,6 +369,8 @@ services ef.UseSqlServer(sqlServerConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.PostgreSql) ef.UsePostgreSql(postgresConnectionString); + else if(sqlDatabaseProvider == SqlDatabaseProvider.MySql) + ef.UseMySql(mySqlConnectionString); else if (sqlDatabaseProvider == SqlDatabaseProvider.CockroachDb) ef.UsePostgreSql(cockroachDbConnectionString!); else @@ -435,6 +447,16 @@ services asb.AzureServiceBusOptions = options => configuration.GetSection("AzureServiceBus").Bind(options); }); } + + if(useKafka) + { + elsa.UseKafka(kafka => + { + 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 new file mode 100644 index 000000000..5acf2f7c6 --- /dev/null +++ b/src/apps/Elsa.Server.Web/WorkflowContextProviders/ConsumerConfigWorkflowContextProvider.cs @@ -0,0 +1,22 @@ +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 new file mode 100644 index 000000000..2f7627367 --- /dev/null +++ b/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs @@ -0,0 +1,33 @@ +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"), + MessageType = new(typeof(ExpandoObject)), + Topics = new(["topic-1"]), + Predicate = new(JavaScriptExpression.Create("message => message.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 new file mode 100644 index 000000000..7f6d0958a --- /dev/null +++ b/src/apps/Elsa.Server.Web/Workflows/ProducerWorkflow.cs @@ -0,0 +1,22 @@ +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 SendMessage + { + 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 90fe55c1b..6db77cdfa 100644 --- a/src/apps/Elsa.Server.Web/appsettings.json +++ b/src/apps/Elsa.Server.Web/appsettings.json @@ -9,6 +9,7 @@ "ConnectionStrings": { "Sqlite": "Data Source=App_Data/elsa.sqlite.db;Cache=Shared;", "PostgreSql": "Server=localhost;Username=elsa;Database=elsa;Port=5432;Password=elsa;SSLMode=Prefer;MaxPoolSize=2000;Timeout=60", + "MySql": "Server=localhost;Port=3306;Database=elsa;User=admin;Password=password;", "CockroachDb": "Host=localhost;Port=26257;Database=elsa;SslMode=Disable;Username=root;IncludeErrorDetail=true", "MongoDb": "mongodb://localhost:27017/elsa-workflows", "AzureServiceBus": "", @@ -196,6 +197,61 @@ } ] }, + "Kafka": { + "Topics": [ + { + "Id": "topic-1", + "Name": "topic-1" + }, + { + "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", + "BootstrapServers": [ + "localhost:9092" + ] + } + ], + "Consumers": [ + { + "Id": "consumer-1", + "Name": "Consumer 1", + "BootstrapServers": [ + "localhost:9092" + ], + "GroupId": "group-1", + "AutoOffsetReset": "earliest", + "EnableAutoCommit": "true", + "CorrelatingFields": [ + "orderId", + "customerId" + ] + }, + { + "Id": "consumer-2", + "Name": "Consumer 2", + "BootstrapServers": [ + "localhost:9092" + ], + "GroupId": "group-1", + "AutoOffsetReset": "earliest", + "EnableAutoCommit": "true" + } + ] + }, "Secrets": { "Management": { "SweepInterval": "04:00:00" diff --git a/src/modules/Elsa.CSharp/Features/CSharpFeature.cs b/src/modules/Elsa.CSharp/Features/CSharpFeature.cs index 9111e5011..48c8b3f96 100644 --- a/src/modules/Elsa.CSharp/Features/CSharpFeature.cs +++ b/src/modules/Elsa.CSharp/Features/CSharpFeature.cs @@ -1,4 +1,5 @@ using Elsa.Common.Features; +using Elsa.CSharp.Activities; using Elsa.CSharp.Contracts; using Elsa.CSharp.Options; using Elsa.CSharp.Providers; @@ -8,6 +9,7 @@ using Elsa.Extensions; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; +using Elsa.Workflows; using Microsoft.Extensions.DependencyInjection; namespace Elsa.CSharp.Features; @@ -45,5 +47,8 @@ public class CSharpFeature : FeatureBase // Activities. Module.AddActivitiesFrom(); + + // UI property handlers. + Services.AddScoped(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Expressions/Models/ExpressionExecutionContext.cs b/src/modules/Elsa.Expressions/Models/ExpressionExecutionContext.cs index 7d19e9adf..1c5ff2806 100644 --- a/src/modules/Elsa.Expressions/Models/ExpressionExecutionContext.cs +++ b/src/modules/Elsa.Expressions/Models/ExpressionExecutionContext.cs @@ -6,49 +6,37 @@ namespace Elsa.Expressions.Models; /// /// Provides context to workflow expressions. /// -public class ExpressionExecutionContext +public class ExpressionExecutionContext( + IServiceProvider serviceProvider, + MemoryRegister memory, + ExpressionExecutionContext? parentContext = default, + IDictionary? transientProperties = default, + CancellationToken cancellationToken = default) { - /// - /// Constructor. - /// - public ExpressionExecutionContext( - IServiceProvider serviceProvider, - MemoryRegister memory, - ExpressionExecutionContext? parentContext = default, - IDictionary? transientProperties = default, - CancellationToken cancellationToken = default) - { - ServiceProvider = serviceProvider; - Memory = memory; - TransientProperties = transientProperties ?? new Dictionary(); - ParentContext = parentContext; - CancellationToken = cancellationToken; - } - /// /// A scoped service provider. /// - public IServiceProvider ServiceProvider { get; } + public IServiceProvider ServiceProvider { get; } = serviceProvider; /// /// A shared register of computer memory. /// - public MemoryRegister Memory { get; } + public MemoryRegister Memory { get; } = memory; /// /// A dictionary of transient properties. /// - public IDictionary TransientProperties { get; set; } + public IDictionary TransientProperties { get; set; } = transientProperties ?? new Dictionary(); /// /// Provides access to the parent , if there is any. /// - public ExpressionExecutionContext? ParentContext { get; set; } + public ExpressionExecutionContext? ParentContext { get; set; } = parentContext; /// /// A cancellation token. /// - public CancellationToken CancellationToken { get; } + public CancellationToken CancellationToken { get; } = cancellationToken; /// /// Returns the pointed to by the specified memory block reference. diff --git a/src/modules/Elsa.Http/Features/HttpFeature.cs b/src/modules/Elsa.Http/Features/HttpFeature.cs index 12bc4e535..3115b1681 100644 --- a/src/modules/Elsa.Http/Features/HttpFeature.cs +++ b/src/modules/Elsa.Http/Features/HttpFeature.cs @@ -186,6 +186,7 @@ public class HttpFeature(IModule module) : FeatureBase(module) // Activity property options providers. .AddScoped() + .AddScoped() .AddScoped(_httpEndpointBasePathProvider) // Port resolvers. diff --git a/src/modules/Elsa.JavaScript/Features/JavaScriptFeature.cs b/src/modules/Elsa.JavaScript/Features/JavaScriptFeature.cs index ee04c38b4..edd70c6bd 100644 --- a/src/modules/Elsa.JavaScript/Features/JavaScriptFeature.cs +++ b/src/modules/Elsa.JavaScript/Features/JavaScriptFeature.cs @@ -36,7 +36,7 @@ public class JavaScriptFeature : FeatureBase /// Configures the Jint options. /// private Action JintOptions { get; set; } = _ => { }; - + public JavaScriptFeature ConfigureJintOptions(Action configure) { JintOptions += configure; @@ -88,8 +88,10 @@ public class JavaScriptFeature : FeatureBase // Activities. Module.UseWorkflowManagement(management => management.AddActivity()); - Services - .AddScoped() - .AddFunctionDefinitionProvider(); + // Type Script definitions. + Services.AddFunctionDefinitionProvider(); + + // UI property handlers. + Services.AddScoped(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Activities/MessageReceived.cs b/src/modules/Elsa.Kafka/Activities/MessageReceived.cs new file mode 100644 index 000000000..c7e455521 --- /dev/null +++ b/src/modules/Elsa.Kafka/Activities/MessageReceived.cs @@ -0,0 +1,130 @@ +using System.Runtime.CompilerServices; +using Elsa.Expressions.Models; +using Elsa.Extensions; +using Elsa.Kafka.Stimuli; +using Elsa.Kafka.UIHints; +using Elsa.Workflows; +using Elsa.Workflows.Attributes; +using Elsa.Workflows.Models; +using Elsa.Workflows.UIHints; +using Microsoft.Extensions.Options; + +namespace Elsa.Kafka.Activities; + +[Activity("Elsa.Kafka", "Kafka", "Executes when a message is received from a given set of topics")] +public class MessageReceived : Trigger +{ + internal const string InputKey = "TransportMessage"; + + /// + public MessageReceived([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) + { + } + + /// + public MessageReceived(Input consumerDefinitionId, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) + { + ConsumerDefinitionId = consumerDefinitionId; + } + + /// + /// The consumer to read from. + /// + [Input( + DisplayName = "Consumer", + Description = "The consumer to connect to.", + UIHandler = typeof(ConsumerDefinitionsDropdownOptionsProvider), + UIHint = InputUIHints.DropDown + )] + public Input ConsumerDefinitionId { get; set; } = default!; + + /// + /// The topics to read from. + /// + [Input( + DisplayName = "Topics", + Description = "The topics to read from.", + UIHint = InputUIHints.MultiText + )] + public Input> Topics { get; set; } = default!; + + /// + /// Optional. The .NET type to deserialize the message into. Defaults to . + /// + [Input(Description = "Optional. The .NET type to deserialize the message into.")] + public Input MessageType { get; set; } = default!; + + [Input( + Description = "Optional. A predicate to filter messages.", + AutoEvaluate = false + )] + public Input Predicate { get; set; } = default!; + + [Input(DisplayName = "Local", Description = "Whether the event is local to the workflow. When checked, only events delivered to this workflow instance will resume this activity.")] + public Input IsLocal { get; set; } = default!; + + /// + /// The received transport message. + /// + public Output TransportMessage = default!; + + /// + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) + { + // If the activity is triggered by a workflow trigger, resume immediately. + if (context.IsTriggerOfWorkflow()) + { + await Resume(context); + } + else + { + // Otherwise, create a bookmark and wait for the stimulus to arrive. + context.CreateBookmark(GetStimulus(context.ExpressionExecutionContext), Resume, false); + } + } + + /// + protected override object GetTriggerPayload(TriggerIndexingContext context) => GetStimulus(context.ExpressionExecutionContext); + + private async ValueTask Resume(ActivityExecutionContext context) + { + var receivedMessage = context.GetWorkflowInput(InputKey); + SetResult(receivedMessage, context); + await context.CompleteActivityAsync(); + } + + private void SetResult(KafkaTransportMessage receivedMessage, ActivityExecutionContext context) + { + var bodyAsString = receivedMessage.Value; + var targetType = MessageType.GetOrDefault(context); + var deserializer = context.GetRequiredService>().Value.Deserializer; + var serviceProvider = context.WorkflowExecutionContext.ServiceProvider; + var body = targetType == null ? bodyAsString : deserializer(serviceProvider, bodyAsString, targetType); + + context.Set(TransportMessage, receivedMessage); + context.SetResult(body); + } + + private object GetStimulus(ExpressionExecutionContext context) + { + var consumerDefinitionId = ConsumerDefinitionId.Get(context); + var topics = Topics.GetOrDefault(context) ?? []; + var messageType = MessageType.GetOrDefault(context) ?? typeof(string); + var isLocal = IsLocal.GetOrDefault(context); + var activity = context.GetActivity(); + var activityRegistry = context.GetRequiredService(); + var activityDescriptor = activityRegistry.Find(activity.Type, activity.Version)!; + var inputDescriptor = activityDescriptor.GetWrappedInputPropertyDescriptor(activity, nameof(Predicate)); + var predicateInput = (Input?)inputDescriptor!.ValueGetter(activity); + var predicateExpression = predicateInput?.Expression; + + return new MessageReceivedStimulus + { + ConsumerDefinitionId = consumerDefinitionId, + Topics = topics.Distinct().ToList(), + MessageType = messageType, + IsLocal = isLocal, + Predicate = predicateExpression, + }; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Activities/SendMessage.cs b/src/modules/Elsa.Kafka/Activities/SendMessage.cs new file mode 100644 index 000000000..9e064e52c --- /dev/null +++ b/src/modules/Elsa.Kafka/Activities/SendMessage.cs @@ -0,0 +1,102 @@ +using System.Text; +using Confluent.Kafka; +using Elsa.Extensions; +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.Options; + +namespace Elsa.Kafka.Activities; + +[Activity("Elsa.Kafka", "Kafka", "Sends a message to a given topic")] +public class SendMessage : CodeActivity +{ + /// + /// The topic to which the message will be sent. + /// + [Input( + Description = "The topic to which the message will be sent.", + UIHint = InputUIHints.DropDown, + UIHandler = typeof(TopicDefinitionsDropdownOptionsProvider) + )] + public Input Topic { get; set; } = default!; + + /// + /// The producer to use when sending the message. + /// + [Input( + DisplayName = "Producer", + Description = "The producer to use when sending the message.", + UIHint = InputUIHints.DropDown, + UIHandler = typeof(ProducerDefinitionsDropdownOptionsProvider) + )] + public Input ProducerDefinitionId { get; set; } = default!; + + [Input(DisplayName = "Local", Description = "When checked, the message will be delivered to this workflow instance only.")] + public Input IsLocal { get; set; } = default!; + + /// + /// Optional. The correlation ID to assign to the message. + /// + [Input( + DisplayName = "Correlation ID", + Description = "Optional. The correlation ID to assign to the message." + )] + public Input CorrelationId { get; set; } = default!; + + /// + /// The content of the message to send. + /// + [Input(Description = "The content of the message to send.")] + public Input Content { get; set; } = default!; + + protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) + { + 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 serializer = context.GetRequiredService>().Value.Serializer; + var serviceProvider = context.WorkflowExecutionContext.ServiceProvider; + var serializedContent = content as string ?? serializer(serviceProvider, content); + var config = new ProducerConfig + { + BootstrapServers = string.Join(",", producerDefinition.BootstrapServers), + }; + + context.DeferTask(async () => + { + using var producer = new ProducerBuilder(config).Build(); + var headers = CreateHeaders(context); + var message = new Message + { + Value = serializedContent, + Headers = headers + }; + await producer.ProduceAsync(topic, message); + }); + } + + private Headers CreateHeaders(ActivityExecutionContext context) + { + var options = context.GetRequiredService>().Value; + var headers = new Headers(); + var correlationId = CorrelationId.GetOrDefault(context); + var isLocal = IsLocal.Get(context); + var topic = Topic.Get(context); + + headers.Add(options.TopicHeaderKey, Encoding.UTF8.GetBytes(topic)); + + if (!string.IsNullOrWhiteSpace(correlationId)) + headers.Add(options.CorrelationHeaderKey, Encoding.UTF8.GetBytes(correlationId)); + + if (isLocal) + headers.Add(options.WorkflowInstanceIdHeaderKey, Encoding.UTF8.GetBytes(context.WorkflowExecutionContext.Id)); + + return headers; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionEnumerator.cs new file mode 100644 index 000000000..6e6964012 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionEnumerator.cs @@ -0,0 +1,11 @@ +namespace Elsa.Kafka; + +public interface IConsumerDefinitionEnumerator +{ + /// + /// Retrieves a list of consumer definitions provided by implementations. + /// + /// A token to monitor for cancellation requests. + /// A list of ConsumerDefinition objects. + Task> EnumerateAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionProvider.cs b/src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionProvider.cs new file mode 100644 index 000000000..f9779d4ce --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionProvider.cs @@ -0,0 +1,6 @@ +namespace Elsa.Kafka; + +public interface IConsumerDefinitionProvider +{ + Task> GetConsumerDefinitionsAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/ICorrelationStrategy.cs b/src/modules/Elsa.Kafka/Contracts/ICorrelationStrategy.cs new file mode 100644 index 000000000..08b794786 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/ICorrelationStrategy.cs @@ -0,0 +1,9 @@ +namespace Elsa.Kafka; + +/// +/// Interface representing a strategy for extracting a correlation ID from a . +/// +public interface ICorrelationStrategy +{ + string? GetCorrelationId(KafkaTransportMessage transportMessage); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/IProducerDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Contracts/IProducerDefinitionEnumerator.cs new file mode 100644 index 000000000..6a8a00d33 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/IProducerDefinitionEnumerator.cs @@ -0,0 +1,16 @@ +namespace Elsa.Kafka; + +public interface IProducerDefinitionEnumerator +{ + /// + /// Retrieves a list of producer definitions provided by implementations. + /// + /// A token to monitor for cancellation requests. + /// A list of ProducerDefinition objects. + Task> EnumerateAsync(CancellationToken cancellationToken); + + /// + /// Retrieves a producer definition by its ID. + /// + Task GetByIdAsync(string id); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/IProducerDefinitionProvider.cs b/src/modules/Elsa.Kafka/Contracts/IProducerDefinitionProvider.cs new file mode 100644 index 000000000..b90b05615 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/IProducerDefinitionProvider.cs @@ -0,0 +1,6 @@ +namespace Elsa.Kafka; + +public interface IProducerDefinitionProvider +{ + Task> GetProducerDefinitionsAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/ITopicDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Contracts/ITopicDefinitionEnumerator.cs new file mode 100644 index 000000000..5db7ec8c6 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/ITopicDefinitionEnumerator.cs @@ -0,0 +1,11 @@ +namespace Elsa.Kafka; + +public interface ITopicDefinitionEnumerator +{ + /// + /// Retrieves a list of topic definitions provided by implementations. + /// + /// A token to monitor for cancellation requests. + /// A list of TopicDefinition objects. + Task> EnumerateAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/ITopicDefinitionProvider.cs b/src/modules/Elsa.Kafka/Contracts/ITopicDefinitionProvider.cs new file mode 100644 index 000000000..bb0ef1d3c --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/ITopicDefinitionProvider.cs @@ -0,0 +1,6 @@ +namespace Elsa.Kafka; + +public interface ITopicDefinitionProvider +{ + Task> GetTopicDefinitionsAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/IWorker.cs b/src/modules/Elsa.Kafka/Contracts/IWorker.cs new file mode 100644 index 000000000..c7cb1f868 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/IWorker.cs @@ -0,0 +1,11 @@ +namespace Elsa.Kafka; + +public interface IWorker : IDisposable +{ + void Start(CancellationToken cancellationToken); + void Stop(); + void BindTrigger(TriggerBinding binding); + void BindBookmark(BookmarkBinding binding); + void RemoveTriggers(IEnumerable triggerIds); + void RemoveBookmarks(IEnumerable bookmarkIds); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/IWorkerManager.cs b/src/modules/Elsa.Kafka/Contracts/IWorkerManager.cs new file mode 100644 index 000000000..e728841fb --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/IWorkerManager.cs @@ -0,0 +1,12 @@ +using Elsa.Workflows.Runtime.Entities; + +namespace Elsa.Kafka; + +public interface IWorkerManager +{ + Task UpdateWorkersAsync(CancellationToken cancellationToken = default); + void StopWorkers(); + Task BindTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default); + Task BindBookmarksAsync(IEnumerable bookmarks, CancellationToken cancellationToken = default); + IWorker? GetWorker(string consumerDefinitionId); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Contracts/IWorkerTopicSubscriber.cs b/src/modules/Elsa.Kafka/Contracts/IWorkerTopicSubscriber.cs new file mode 100644 index 000000000..99bd1a104 --- /dev/null +++ b/src/modules/Elsa.Kafka/Contracts/IWorkerTopicSubscriber.cs @@ -0,0 +1,9 @@ +namespace Elsa.Kafka; + +public interface IWorkerTopicSubscriber +{ + /// + /// Updates all workers by subscribing to topics based on the current workflow trigger indexes and bookmarks. + /// + Task UpdateTopicSubscriptionsAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Correlation/HeaderCorrelationStrategy.cs b/src/modules/Elsa.Kafka/Correlation/HeaderCorrelationStrategy.cs new file mode 100644 index 000000000..f5e88ffbb --- /dev/null +++ b/src/modules/Elsa.Kafka/Correlation/HeaderCorrelationStrategy.cs @@ -0,0 +1,19 @@ +using System.Text; +using Microsoft.Extensions.Options; + +namespace Elsa.Kafka; + +/// +/// A correlation strategy that retrieves a correlation ID from a specified header in a Kafka transport message. +/// +public class HeaderCorrelationStrategy(IOptions options) : ICorrelationStrategy +{ + /// + public string? GetCorrelationId(KafkaTransportMessage transportMessage) + { + var headerName = options.Value.CorrelationHeaderKey; + return !transportMessage.Headers.TryGetValue(headerName, out var headerValue) + ? null + : Encoding.UTF8.GetString(headerValue); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Correlation/NullCorrelationStrategy.cs b/src/modules/Elsa.Kafka/Correlation/NullCorrelationStrategy.cs new file mode 100644 index 000000000..9788655bd --- /dev/null +++ b/src/modules/Elsa.Kafka/Correlation/NullCorrelationStrategy.cs @@ -0,0 +1,11 @@ +namespace Elsa.Kafka; + +/// +/// Represents a strategy that does not extract any correlation ID from a . +/// This implementation of always returns null, +/// effectively indicating no correlation ID is available or applicable. +/// +public class NullCorrelationStrategy : ICorrelationStrategy +{ + public string? GetCorrelationId(KafkaTransportMessage transportMessage) => null; +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Elsa.Kafka.csproj b/src/modules/Elsa.Kafka/Elsa.Kafka.csproj new file mode 100644 index 000000000..fb8e823ad --- /dev/null +++ b/src/modules/Elsa.Kafka/Elsa.Kafka.csproj @@ -0,0 +1,24 @@ + + + + + Provides Apache Kafka integration and activities. + + elsa module kafka service-bus + + + + + + + + + + + + + + + + + diff --git a/src/modules/Elsa.Kafka/Elsa.Kafka.csproj.DotSettings b/src/modules/Elsa.Kafka/Elsa.Kafka.csproj.DotSettings new file mode 100644 index 000000000..ed333bc2b --- /dev/null +++ b/src/modules/Elsa.Kafka/Elsa.Kafka.csproj.DotSettings @@ -0,0 +1,9 @@ + + True + True + True + True + True + True + True + \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Entities/ConsumerConfigDefinition.cs b/src/modules/Elsa.Kafka/Entities/ConsumerConfigDefinition.cs new file mode 100644 index 000000000..8de9dc460 --- /dev/null +++ b/src/modules/Elsa.Kafka/Entities/ConsumerConfigDefinition.cs @@ -0,0 +1,13 @@ +using Confluent.Kafka; +using Elsa.Common.Entities; + +namespace Elsa.Kafka; + +public class ConsumerDefinition : Entity +{ + public string Name { get; set; } = default!; + public ICollection BootstrapServers { get; set; } = []; + public string GroupId { get; set; } = default!; + public AutoOffsetReset AutoOffsetReset { get; set; } = AutoOffsetReset.Earliest; + public bool EnableAutoCommit { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs b/src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs new file mode 100644 index 000000000..f270b4773 --- /dev/null +++ b/src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs @@ -0,0 +1,9 @@ +using Elsa.Common.Entities; + +namespace Elsa.Kafka; + +public class ProducerDefinition : Entity +{ + public string Name { get; set; } = default!; + public ICollection BootstrapServers { get; set; } = []; +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Entities/TopicDefinition.cs b/src/modules/Elsa.Kafka/Entities/TopicDefinition.cs new file mode 100644 index 000000000..861ddd5e8 --- /dev/null +++ b/src/modules/Elsa.Kafka/Entities/TopicDefinition.cs @@ -0,0 +1,11 @@ +using Elsa.Common.Entities; + +namespace Elsa.Kafka; + +public class TopicDefinition : Entity +{ + /// + /// Gets or sets the name of the topic. + /// + public string Name { get; set; } = default!; +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Kafka/Extensions/ModuleExtensions.cs new file mode 100644 index 000000000..c8e652e0d --- /dev/null +++ b/src/modules/Elsa.Kafka/Extensions/ModuleExtensions.cs @@ -0,0 +1,12 @@ +using Elsa.Extensions; +using Elsa.Features.Services; + +namespace Elsa.Kafka; + +public static class ModuleExtensions +{ + public static IModule UseKafka(this IModule module, Action? configure = null) + { + return module.Use(configure); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Features/KafkaFeature.cs b/src/modules/Elsa.Kafka/Features/KafkaFeature.cs new file mode 100644 index 000000000..e8ce7fd98 --- /dev/null +++ b/src/modules/Elsa.Kafka/Features/KafkaFeature.cs @@ -0,0 +1,76 @@ +using Elsa.Extensions; +using Elsa.Features.Abstractions; +using Elsa.Features.Services; +using Elsa.Kafka.Implementations; +using Elsa.Kafka.Providers; +using Elsa.Kafka.Tasks; +using Elsa.Kafka.UIHints; +using Elsa.Workflows; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Kafka; + +public class KafkaFeature(IModule module) : FeatureBase(module) +{ + private Action _configureOptions = _ => { }; + private Func _correlationStrategyFactory = sp => sp.GetRequiredService(); + + public KafkaFeature ConfigureOptions(Action configureOptions) + { + _configureOptions += configureOptions; + return this; + } + + public KafkaFeature WithCorrelationStrategy() where T : class, ICorrelationStrategy + { + Services.AddScoped(); + _correlationStrategyFactory = sp => sp.GetRequiredService(); + return this; + } + + public KafkaFeature WithCorrelationStrategy(Func correlationStrategyFactory) + { + _correlationStrategyFactory = correlationStrategyFactory; + return this; + } + + public KafkaFeature WithCorrelationStrategy(Func correlationStrategyFactory) + { + _correlationStrategyFactory = _ => correlationStrategyFactory(); + return this; + } + + public KafkaFeature WithCorrelationStrategy(ICorrelationStrategy correlationStrategy) + { + _correlationStrategyFactory = _ => correlationStrategy; + return this; + } + + public override void Configure() + { + Module.AddActivitiesFrom(); + } + + public override void Apply() + { + Services.Configure(_configureOptions); + + Services + .AddBackgroundTask() + .AddSingleton() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped() + .AddScoped(_correlationStrategyFactory) + .AddHandlersFrom(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/FodyWeavers.xml b/src/modules/Elsa.Kafka/FodyWeavers.xml new file mode 100644 index 000000000..00e1d9a1c --- /dev/null +++ b/src/modules/Elsa.Kafka/FodyWeavers.xml @@ -0,0 +1,3 @@ + + + \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs b/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs new file mode 100644 index 000000000..f04f12335 --- /dev/null +++ b/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs @@ -0,0 +1,208 @@ +using System.Text; +using Elsa.Expressions.Contracts; +using Elsa.Expressions.Models; +using Elsa.Kafka.Activities; +using Elsa.Kafka.Notifications; +using Elsa.Kafka.Stimuli; +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Helpers; +using Elsa.Workflows.Memory; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Options; +using JetBrains.Annotations; +using Microsoft.Extensions.Options; + +namespace Elsa.Kafka.Handlers; + +[UsedImplicitly] +public class TriggerWorkflows( + ITriggerInvoker triggerInvoker, + IBookmarkQueue bookmarkQueue, + ICorrelationStrategy correlationStrategy, + IExpressionEvaluator expressionEvaluator, + IOptions options, + IServiceProvider serviceProvider) : INotificationHandler +{ + private static readonly string MessageReceivedActivityTypeName = ActivityTypeNameHelper.GenerateTypeName(); + + public async Task HandleAsync(TransportMessageReceived notification, CancellationToken cancellationToken) + { + var worker = notification.Worker; + var consumerDefinition = worker.ConsumerDefinition; + var consumerDefinitionId = consumerDefinition.Id; + var transportMessage = notification.TransportMessage; + var boundTriggers = worker.TriggerBindings.Values; + var boundBookmarks = worker.BookmarkBindings.Values; + var matchingTriggers = await GetMatchingTriggerBindingsAsync(boundTriggers, consumerDefinitionId, transportMessage, cancellationToken); + var matchingBookmarks = await GetMatchingBookmarkBindingsAsync(boundBookmarks, consumerDefinitionId, transportMessage, cancellationToken); + + await InvokeTriggersAsync(matchingTriggers, transportMessage, cancellationToken); + await InvokeBookmarksAsync(matchingBookmarks, transportMessage, cancellationToken); + } + + private async Task InvokeTriggersAsync(IEnumerable matchingTriggers, KafkaTransportMessage transportMessage, CancellationToken cancellationToken) + { + var input = new Dictionary + { + [MessageReceived.InputKey] = transportMessage + }; + + foreach (var binding in matchingTriggers) + { + var invokeTriggerRequest = new InvokeTriggerRequest + { + Workflow = binding.Workflow, + ActivityId = binding.TriggerActivityId, + CorrelationId = GetCorrelationId(transportMessage), + Input = input, + Properties = new Dictionary + { + [MessageReceived.InputKey] = transportMessage + } + }; + await triggerInvoker.InvokeAsync(invokeTriggerRequest, cancellationToken); + } + } + + private async Task InvokeBookmarksAsync(IEnumerable matchingBookmarks, KafkaTransportMessage transportMessage, CancellationToken cancellationToken) + { + var input = new Dictionary + { + [MessageReceived.InputKey] = transportMessage + }; + + var properties = new Dictionary + { + [MessageReceived.InputKey] = transportMessage + }; + + foreach (var binding in matchingBookmarks) + { + var bookmarkQueueItem = new NewBookmarkQueueItem + { + WorkflowInstanceId = binding.WorkflowInstanceId, + BookmarkId = binding.BookmarkId, + Options = new ResumeBookmarkOptions + { + Input = input, + Properties = properties + }, + ActivityTypeName = MessageReceivedActivityTypeName + }; + await bookmarkQueue.EnqueueAsync(bookmarkQueueItem, cancellationToken); + } + } + + private async Task> GetMatchingTriggerBindingsAsync( + IEnumerable boundTriggers, + string consumerDefinitionId, + KafkaTransportMessage transportMessage, + CancellationToken cancellationToken) + { + var matchingTriggers = new List(); + var topic = GetTopic(transportMessage); + + if (string.IsNullOrEmpty(topic)) + return matchingTriggers; + + foreach (var binding in boundTriggers) + { + var stimulus = binding.Stimulus; + if (stimulus.ConsumerDefinitionId != consumerDefinitionId) + continue; + + if (stimulus.Topics.All(x => x != topic)) + continue; + + var isMatch = await EvaluatePredicateAsync(transportMessage, stimulus, cancellationToken); + + if (isMatch) + matchingTriggers.Add(binding); + } + + return matchingTriggers; + } + + private async Task> GetMatchingBookmarkBindingsAsync( + IEnumerable boundBookmarks, + string consumerDefinitionId, + KafkaTransportMessage transportMessage, + CancellationToken cancellationToken) + { + var matchingBookmarks = new List(); + var topic = GetTopic(transportMessage); + + if (string.IsNullOrEmpty(topic)) + return matchingBookmarks; + + var correlationId = GetCorrelationId(transportMessage); + var workflowInstanceId = GetWorkflowInstanceId(transportMessage); + + foreach (var binding in boundBookmarks) + { + var stimulus = binding.Stimulus; + if (stimulus.ConsumerDefinitionId != consumerDefinitionId) + continue; + + if (stimulus.Topics.All(x => x != topic)) + continue; + + if (stimulus.IsLocal) + { + if (binding.WorkflowInstanceId != workflowInstanceId) + { + if (string.IsNullOrWhiteSpace(correlationId)) + continue; + + if (binding.CorrelationId != correlationId) + continue; + } + } + + var isMatch = await EvaluatePredicateAsync(transportMessage, stimulus, cancellationToken); + + if (isMatch) + matchingBookmarks.Add(binding); + } + + return matchingBookmarks; + } + + private async Task EvaluatePredicateAsync(KafkaTransportMessage transportMessage, MessageReceivedStimulus stimulus, CancellationToken cancellationToken) + { + var predicate = stimulus.Predicate; + + if (predicate == null) + return true; + + var memory = new MemoryRegister(); + var messageVariable = new Variable("message", transportMessage); + var messageType = stimulus.MessageType; + var message = DeserializeMessage(transportMessage.Value, messageType); + var expressionExecutionContext = new ExpressionExecutionContext(serviceProvider, memory, cancellationToken: cancellationToken); + messageVariable.Set(expressionExecutionContext, message); + return await expressionEvaluator.EvaluateAsync(predicate, expressionExecutionContext); + } + + private object DeserializeMessage(string value, Type? type) + { + return type == null ? value : options.Value.Deserializer(serviceProvider, value, type); + } + + private string? GetWorkflowInstanceId(KafkaTransportMessage transportMessage) + { + var key = options.Value.WorkflowInstanceIdHeaderKey; + return transportMessage.Headers.TryGetValue(key, out var value) ? Encoding.UTF8.GetString(value) : null; + } + + private string? GetCorrelationId(KafkaTransportMessage transportMessage) + { + return correlationStrategy.GetCorrelationId(transportMessage); + } + + private string? GetTopic(KafkaTransportMessage transportMessage) + { + var key = options.Value.TopicHeaderKey; + return transportMessage.Headers.TryGetValue(key, out var value) ? Encoding.UTF8.GetString(value) : null; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Handlers/UpdateWorkers.cs b/src/modules/Elsa.Kafka/Handlers/UpdateWorkers.cs new file mode 100644 index 000000000..bfab0e070 --- /dev/null +++ b/src/modules/Elsa.Kafka/Handlers/UpdateWorkers.cs @@ -0,0 +1,102 @@ +using Elsa.Extensions; +using Elsa.Kafka.Activities; +using Elsa.Kafka.Stimuli; +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Helpers; +using Elsa.Workflows.Management; +using Elsa.Workflows.Models; +using Elsa.Workflows.Runtime.Notifications; +using JetBrains.Annotations; + +namespace Elsa.Kafka.Handlers; + +/// +/// Creates workers for each trigger & bookmark in response to updated workflow trigger indexes and bookmarks. +/// +[UsedImplicitly] +public class UpdateWorkers(IWorkerManager workerManager, IWorkflowDefinitionService workflowDefinitionService) : + INotificationHandler, + INotificationHandler, + INotificationHandler +{ + private static readonly string MessageReceivedActivityTypeName = ActivityTypeNameHelper.GenerateTypeName(); + + /// Adds, updates and removes workers based on added and removed triggers. + public async Task HandleAsync(WorkflowTriggersIndexed notification, CancellationToken cancellationToken) + { + var removedTriggers = notification.IndexedWorkflowTriggers.RemovedTriggers.Where(x => x.Name == MessageReceivedActivityTypeName).ToList(); + var removedTriggerIds = removedTriggers.Select(x => x.Id).ToList(); + + foreach (var removedTrigger in removedTriggers) + { + var consumerDefinitionId = removedTrigger.GetPayload().ConsumerDefinitionId; + var worker = workerManager.GetWorker(consumerDefinitionId); + + worker?.RemoveTriggers(removedTriggerIds); + } + + var addedTriggers = notification.IndexedWorkflowTriggers.AddedTriggers.Where(x => x.Name == MessageReceivedActivityTypeName).ToList(); + + foreach (var addedTrigger in addedTriggers) + { + var stimulus = addedTrigger.GetPayload(); + var consumerDefinitionId = stimulus.ConsumerDefinitionId; + var worker = workerManager.GetWorker(consumerDefinitionId); + + if (worker == null) + continue; + + var workflow = await workflowDefinitionService.FindWorkflowAsync(addedTrigger.WorkflowDefinitionVersionId, cancellationToken); + + if (workflow == null) + continue; + + var triggerBinding = new TriggerBinding(workflow, addedTrigger.Id, addedTrigger.ActivityId, stimulus); + worker.BindTrigger(triggerBinding); + } + } + + /// Adds, updates and removes workers based on added and removed bookmarks. + public Task HandleAsync(WorkflowBookmarksIndexed notification, CancellationToken cancellationToken) + { + RemoveBookmarkBindings(notification.IndexedWorkflowBookmarks.RemovedBookmarks.Where(x => x.Name == MessageReceivedActivityTypeName).ToList()); + + var addedBookmarks = notification.IndexedWorkflowBookmarks.AddedBookmarks.Where(x => x.Name == MessageReceivedActivityTypeName).ToList(); + var workflowInstanceId = notification.IndexedWorkflowBookmarks.WorkflowExecutionContext.Id; + var correlationId = notification.IndexedWorkflowBookmarks.WorkflowExecutionContext.CorrelationId; + + foreach (var addedBookmark in addedBookmarks) + { + var stimulus = addedBookmark.GetPayload(); + var consumerDefinitionId = stimulus.ConsumerDefinitionId; + var worker = workerManager.GetWorker(consumerDefinitionId); + + if (worker == null) + continue; + + var bookmarkBinding = new BookmarkBinding(workflowInstanceId, correlationId, addedBookmark.Id, stimulus); + worker.BindBookmark(bookmarkBinding); + } + + return Task.CompletedTask; + } + + public Task HandleAsync(BookmarksDeleted notification, CancellationToken cancellationToken) + { + RemoveBookmarkBindings(notification.Bookmarks.Where(x => x.ActivityTypeName == MessageReceivedActivityTypeName).Select(x => x.ToBookmark()).ToList()); + return Task.CompletedTask; + } + + private void RemoveBookmarkBindings(ICollection bookmarks) + { + var removedBookmarkIds = bookmarks.Select(x => x.Id).ToList(); + foreach (var bookmark in bookmarks) + { + var stimulus = bookmark.GetPayload(); + var consumerDefinitionId = stimulus.ConsumerDefinitionId; + var worker = workerManager.GetWorker(consumerDefinitionId); + + worker?.RemoveBookmarks(removedBookmarkIds); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Implementations/ConsumerDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Implementations/ConsumerDefinitionEnumerator.cs new file mode 100644 index 000000000..4683202e4 --- /dev/null +++ b/src/modules/Elsa.Kafka/Implementations/ConsumerDefinitionEnumerator.cs @@ -0,0 +1,23 @@ +using System.Runtime.CompilerServices; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.Implementations; + +public class ConsumerDefinitionEnumerator(IEnumerable providers) : IConsumerDefinitionEnumerator +{ + public async Task> EnumerateAsync(CancellationToken cancellationToken) + { + return await GetConsumerDefinitionsInternalAsync(cancellationToken).ToListAsync(cancellationToken); + } + + private async IAsyncEnumerable GetConsumerDefinitionsInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + foreach (var provider in providers) + { + var consumerDefinitions = await provider.GetConsumerDefinitionsAsync(cancellationToken).ToList(); + + foreach (var consumerDefinition in consumerDefinitions) + yield return consumerDefinition; + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Implementations/ProducerDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Implementations/ProducerDefinitionEnumerator.cs new file mode 100644 index 000000000..bf2ad9b7c --- /dev/null +++ b/src/modules/Elsa.Kafka/Implementations/ProducerDefinitionEnumerator.cs @@ -0,0 +1,29 @@ +using System.Runtime.CompilerServices; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.Implementations; + +public class ProducerDefinitionEnumerator(IEnumerable providers) : IProducerDefinitionEnumerator +{ + public async Task> EnumerateAsync(CancellationToken cancellationToken) + { + return await GetProducerDefinitionsInternalAsync(cancellationToken).ToListAsync(cancellationToken); + } + + public async Task GetByIdAsync(string id) + { + var producerDefinitions = await GetProducerDefinitionsInternalAsync(CancellationToken.None).ToListAsync(); + return producerDefinitions.FirstOrDefault(x => x.Id == id) ?? throw new InvalidOperationException($"Producer definition with ID '{id}' not found."); + } + + private async IAsyncEnumerable GetProducerDefinitionsInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + foreach (var provider in providers) + { + var producerDefinitions = await provider.GetProducerDefinitionsAsync(cancellationToken).ToList(); + + foreach (var producerDefinition in producerDefinitions) + yield return producerDefinition; + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Implementations/TopicDefinitionEnumerator.cs b/src/modules/Elsa.Kafka/Implementations/TopicDefinitionEnumerator.cs new file mode 100644 index 000000000..fe972ce9a --- /dev/null +++ b/src/modules/Elsa.Kafka/Implementations/TopicDefinitionEnumerator.cs @@ -0,0 +1,23 @@ +using System.Runtime.CompilerServices; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.Implementations; + +public class TopicDefinitionEnumerator(IEnumerable providers) : ITopicDefinitionEnumerator +{ + public async Task> EnumerateAsync(CancellationToken cancellationToken) + { + return await GetTopicDefinitionsInternalAsync(cancellationToken).ToListAsync(cancellationToken); + } + + private async IAsyncEnumerable GetTopicDefinitionsInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken) + { + foreach (var provider in providers) + { + var topicDefinitions = await provider.GetTopicDefinitionsAsync(cancellationToken).ToList(); + + foreach (var topicDefinition in topicDefinitions) + yield return topicDefinition; + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Implementations/Worker.cs b/src/modules/Elsa.Kafka/Implementations/Worker.cs new file mode 100644 index 000000000..42b3b1c26 --- /dev/null +++ b/src/modules/Elsa.Kafka/Implementations/Worker.cs @@ -0,0 +1,128 @@ +using Confluent.Kafka; +using Elsa.Extensions; + +namespace Elsa.Kafka.Implementations; + +public class Worker(ConsumerDefinition consumerDefinition) : IWorker +{ + private bool _running; + private IConsumer _consumer = default!; + private CancellationTokenSource _cancellationTokenSource = new(); + private readonly HashSet _subscribedTopics = new(); + + public IDictionary BookmarkBindings { get; } = new Dictionary(); + public IDictionary TriggerBindings { get; } = new Dictionary(); + public Func, CancellationToken, Task>? MessageReceived { get; set; } + public ConsumerDefinition ConsumerDefinition { get; } = consumerDefinition; + + public void Start(CancellationToken cancellationToken) + { + if (_running) + return; + + _running = true; + + var config = new ConsumerConfig + { + BootstrapServers = string.Join(",", ConsumerDefinition.BootstrapServers), + GroupId = ConsumerDefinition.GroupId, + AutoOffsetReset = ConsumerDefinition.AutoOffsetReset, + EnableAutoCommit = ConsumerDefinition.EnableAutoCommit + }; + + _consumer = new ConsumerBuilder(config).Build(); + _cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + _ = Task.Run(() => RunAsync(_cancellationTokenSource.Token), cancellationToken); + } + + public void Stop() + { + if (!_running) + return; + + _running = false; + _consumer.Unsubscribe(); + _consumer.Close(); + _consumer.Dispose(); + _cancellationTokenSource.Cancel(); + _cancellationTokenSource.Dispose(); + } + + public void BindTrigger(TriggerBinding binding) + { + var topics = binding.Stimulus.Topics.Distinct().ToList(); + TriggerBindings[binding.TriggerId] = binding; + Subscribe(topics); + } + + public void BindBookmark(BookmarkBinding binding) + { + var topics = binding.Stimulus.Topics.Distinct().ToList(); + BookmarkBindings[binding.BookmarkId] = binding; + Subscribe(topics); + } + + public void RemoveTriggers(IEnumerable triggerIds) + { + var triggerIdList = triggerIds.ToList(); + TriggerBindings.RemoveWhere(x => triggerIdList.Contains(x.Key)); + } + + public void RemoveBookmarks(IEnumerable bookmarkIds) + { + var bookmarkIdList = bookmarkIds.ToList(); + BookmarkBindings.RemoveWhere(x => bookmarkIdList.Contains(x.Key)); + } + + public void Subscribe(IEnumerable bookmarkBindings) + { + var topics = bookmarkBindings.SelectMany(x => x.Stimulus.Topics).Distinct().ToList(); + Subscribe(topics); + } + + public void Dispose() + { + Stop(); + } + + private void Subscribe(IEnumerable topics) + { + if(!_running) + throw new InvalidOperationException("The worker is not running."); + + var topicList = topics.ToHashSet(); + + // Check if there are any topics not yet in the list of subscribed topics. + if (topicList.All(topic => _subscribedTopics.Contains(topic))) + return; + + // Add the new topics to the list of subscribed topics. + _subscribedTopics.UnionWith(topicList); + + // Update the consumer's subscription. + _consumer.Subscribe(_subscribedTopics); + } + + private async Task RunAsync(CancellationToken cancellationToken) + { + while (!cancellationToken.IsCancellationRequested) + { + var consumeResult = _consumer.Consume(cancellationToken); + + if (consumeResult.IsPartitionEOF) + continue; + + await ProcessMessageAsync(consumeResult.Message); + } + + _consumer.Unsubscribe(); + _consumer.Close(); + _consumer.Dispose(); + } + + private Task ProcessMessageAsync(Message message) + { + MessageReceived?.Invoke(this, message, _cancellationTokenSource.Token); + return Task.CompletedTask; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Implementations/WorkerManager.cs b/src/modules/Elsa.Kafka/Implementations/WorkerManager.cs new file mode 100644 index 000000000..3f9276e83 --- /dev/null +++ b/src/modules/Elsa.Kafka/Implementations/WorkerManager.cs @@ -0,0 +1,136 @@ +using Confluent.Kafka; +using Elsa.Extensions; +using Elsa.Kafka.Notifications; +using Elsa.Kafka.Stimuli; +using Elsa.Mediator.Contracts; +using Elsa.Workflows; +using Elsa.Workflows.Management; +using Elsa.Workflows.Runtime.Entities; +using Microsoft.Extensions.DependencyInjection; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.Implementations; + +public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) : IWorkerManager +{ + private IDictionary Workers { get; set; } = new Dictionary(); + + public async Task UpdateWorkersAsync(CancellationToken cancellationToken = default) + { + await using var scope = scopeFactory.CreateAsyncScope(); + var consumerDefinitionEnumerator = scope.ServiceProvider.GetRequiredService(); + var consumerDefinitions = await consumerDefinitionEnumerator.EnumerateAsync(cancellationToken).ToList(); + var workers = new Dictionary(Workers); + var workersToRemove = workers.Keys.Except(consumerDefinitions.Select(x => x.Id)).ToList(); + + // Remove workers that are no longer needed. + foreach (var workerToRemove in workersToRemove) + { + if (workers.TryGetValue(workerToRemove, out var worker)) + { + worker.Stop(); + workers.Remove(workerToRemove); + } + } + + // Add or update workers. + foreach (var consumerDefinition in consumerDefinitions) + { + var existingWorker = workers.GetValueOrDefault(consumerDefinition.Id); + + if (existingWorker == null) + { + // Add a new worker. + var worker = CreateWorker(consumerDefinition); + workers.Add(consumerDefinition.Id, worker); + } + else + { + // Compare the existing worker's consumer definition with the new consumer definition using a hash. + var existingHash = ComputeHash(existingWorker.ConsumerDefinition); + var hash = ComputeHash(consumerDefinition); + + // If the hash is different, update the worker. + if (existingHash != hash) + { + existingWorker.Stop(); + var worker = CreateWorker(consumerDefinition); + workers[consumerDefinition.Id] = worker; + } + } + } + + // Ensure all workers are running. + foreach (var worker in workers.Values) + worker.Start(cancellationToken); + + // Update the worker dictionary. + Workers = workers; + } + + public async Task BindTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default) + { + await using var scope = scopeFactory.CreateAsyncScope(); + var workflowDefinitionService = scope.ServiceProvider.GetRequiredService(); + + foreach (var trigger in triggers) + { + var workflow = await workflowDefinitionService.FindWorkflowAsync(trigger.WorkflowDefinitionVersionId, cancellationToken); + + if (workflow == null) + continue; + + var stimulus = trigger.GetPayload(); + var consumerDefinitionId = stimulus.ConsumerDefinitionId; + var worker = GetWorker(consumerDefinitionId); + var triggerBinding = new TriggerBinding(workflow, trigger.Id, trigger.ActivityId, stimulus); + worker.BindTrigger(triggerBinding); + } + } + + public Task BindBookmarksAsync(IEnumerable bookmarks, CancellationToken cancellationToken = default) + { + foreach (var bookmark in bookmarks) + { + var stimulus = bookmark.GetPayload(); + var consumerDefinitionId = stimulus.ConsumerDefinitionId; + var worker = GetWorker(consumerDefinitionId); + var bookmarkBinding = new BookmarkBinding(bookmark.WorkflowInstanceId, bookmark.CorrelationId, bookmark.Id, stimulus); + worker.BindBookmark(bookmarkBinding); + } + + return Task.CompletedTask; + } + + public void StopWorkers() + { + foreach (var worker in Workers.Values) + worker.Stop(); + } + + public IWorker? GetWorker(string consumerDefinitionId) + { + return Workers.TryGetValue(consumerDefinitionId, out var worker) ? worker : null; + } + + private Worker CreateWorker(ConsumerDefinition consumerDefinition) + { + var worker = new Worker(consumerDefinition); + worker.MessageReceived = OnMessageReceivedAsync; + return worker; + } + + private async Task OnMessageReceivedAsync(Worker worker, Message arg, CancellationToken cancellationToken) + { + var headers = arg.Headers.ToDictionary(x => x.Key, x => x.GetValueBytes()); + var notification = new TransportMessageReceived(worker, new KafkaTransportMessage(arg.Key, arg.Value, headers)); + await using var scope = scopeFactory.CreateAsyncScope(); + var mediator = scope.ServiceProvider.GetRequiredService(); + await mediator.SendAsync(notification, cancellationToken); + } + + private string ComputeHash(ConsumerDefinition consumerDefinition) + { + return hasher.Hash(consumerDefinition); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Implementations/WorkerTopicSubscriber.cs b/src/modules/Elsa.Kafka/Implementations/WorkerTopicSubscriber.cs new file mode 100644 index 000000000..bc5fecdb4 --- /dev/null +++ b/src/modules/Elsa.Kafka/Implementations/WorkerTopicSubscriber.cs @@ -0,0 +1,38 @@ +using Elsa.Kafka.Activities; +using Elsa.Workflows.Helpers; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Entities; +using Elsa.Workflows.Runtime.Filters; + +namespace Elsa.Kafka.Implementations; + +public class WorkerTopicSubscriber(ITriggerStore triggerStore, IBookmarkStore bookmarkStore, IWorkerManager workerManager) : IWorkerTopicSubscriber +{ + private static readonly string MessageReceivedActivityTypeName = ActivityTypeNameHelper.GenerateTypeName(); + + public async Task UpdateTopicSubscriptionsAsync(CancellationToken cancellationToken) + { + var triggers = await GetStoredTriggersAsync(cancellationToken); + var bookmarks = await GetStoredBookmarksAsync(cancellationToken); + await workerManager.BindTriggersAsync(triggers, cancellationToken); + await workerManager.BindBookmarksAsync(bookmarks, cancellationToken); + } + + private async Task> GetStoredTriggersAsync(CancellationToken cancellationToken) + { + var triggerFilter = new TriggerFilter + { + Name = MessageReceivedActivityTypeName + }; + return await triggerStore.FindManyAsync(triggerFilter, cancellationToken); + } + + private async Task> GetStoredBookmarksAsync(CancellationToken cancellationToken) + { + var bookmarkFilter = new BookmarkFilter + { + ActivityTypeName = MessageReceivedActivityTypeName + }; + return await bookmarkStore.FindManyAsync(bookmarkFilter, cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Models/BookmarkBinding.cs b/src/modules/Elsa.Kafka/Models/BookmarkBinding.cs new file mode 100644 index 000000000..1a310655b --- /dev/null +++ b/src/modules/Elsa.Kafka/Models/BookmarkBinding.cs @@ -0,0 +1,5 @@ +using Elsa.Kafka.Stimuli; + +namespace Elsa.Kafka; + +public record BookmarkBinding(string WorkflowInstanceId, string? CorrelationId, string BookmarkId, MessageReceivedStimulus Stimulus); \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Models/KafkaTransportMessage.cs b/src/modules/Elsa.Kafka/Models/KafkaTransportMessage.cs new file mode 100644 index 000000000..1dfb58fac --- /dev/null +++ b/src/modules/Elsa.Kafka/Models/KafkaTransportMessage.cs @@ -0,0 +1,3 @@ +namespace Elsa.Kafka; + +public record KafkaTransportMessage(object? Key, string Value, IDictionary Headers); \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Models/TriggerBinding.cs b/src/modules/Elsa.Kafka/Models/TriggerBinding.cs new file mode 100644 index 000000000..d6368fc64 --- /dev/null +++ b/src/modules/Elsa.Kafka/Models/TriggerBinding.cs @@ -0,0 +1,6 @@ +using Elsa.Kafka.Stimuli; +using Elsa.Workflows.Activities; + +namespace Elsa.Kafka; + +public record TriggerBinding(Workflow Workflow, string TriggerId, string TriggerActivityId, MessageReceivedStimulus Stimulus); \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Notifications/TransportMessageReceived.cs b/src/modules/Elsa.Kafka/Notifications/TransportMessageReceived.cs new file mode 100644 index 000000000..259e30dcd --- /dev/null +++ b/src/modules/Elsa.Kafka/Notifications/TransportMessageReceived.cs @@ -0,0 +1,6 @@ +using Elsa.Kafka.Implementations; +using Elsa.Mediator.Contracts; + +namespace Elsa.Kafka.Notifications; + +public record TransportMessageReceived(Worker Worker, KafkaTransportMessage TransportMessage) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Options/KafkaOptions.cs b/src/modules/Elsa.Kafka/Options/KafkaOptions.cs new file mode 100644 index 000000000..c29b5403a --- /dev/null +++ b/src/modules/Elsa.Kafka/Options/KafkaOptions.cs @@ -0,0 +1,16 @@ +using Elsa.Kafka.Serialization; + +namespace Elsa.Kafka; + +public class KafkaOptions +{ + public ICollection Topics { get; set; } = []; + public ICollection Consumers { get; set; } = []; + public ICollection Producers { get; set; } = []; + public string WorkflowInstanceIdHeaderKey { get; set; } = "x-workflow-instance-id"; + public string CorrelationHeaderKey { get; set; } = "x-correlation-id"; + public string TopicHeaderKey { get; set; } = "x-topic"; + + public Func Serializer { get; set; } = DefaultSerializers.SerializePayload; + public Func Deserializer { get; set; } = DefaultSerializers.DeserializePayload; +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs b/src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs new file mode 100644 index 000000000..668146b5b --- /dev/null +++ b/src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs @@ -0,0 +1,21 @@ +using Microsoft.Extensions.Options; + +namespace Elsa.Kafka.Providers; + +public class OptionsDefinitionProvider(IOptions options) : IConsumerDefinitionProvider, IProducerDefinitionProvider, ITopicDefinitionProvider +{ + public Task> GetConsumerDefinitionsAsync(CancellationToken cancellationToken = default) + { + return Task.FromResult>(options.Value.Consumers); + } + + public Task> GetProducerDefinitionsAsync(CancellationToken cancellationToken = default) + { + return Task.FromResult>(options.Value.Producers); + } + + public Task> GetTopicDefinitionsAsync(CancellationToken cancellationToken = default) + { + return Task.FromResult>(options.Value.Topics); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Serialization/DefaultSerializers.cs b/src/modules/Elsa.Kafka/Serialization/DefaultSerializers.cs new file mode 100644 index 000000000..8598826cc --- /dev/null +++ b/src/modules/Elsa.Kafka/Serialization/DefaultSerializers.cs @@ -0,0 +1,35 @@ +using System.Text.Json; +using Elsa.Workflows.Serialization.Converters; + +namespace Elsa.Kafka.Serialization; + +public static class DefaultSerializers +{ + private static readonly JsonSerializerOptions SerializerOptions = GetOptions(); + + private static JsonSerializerOptions GetOptions() + { + var options = new JsonSerializerOptions + { + PropertyNamingPolicy = JsonNamingPolicy.CamelCase, + Converters = + { + new ExpandoObjectConverterFactory() + } + }; + + return options; + } + + public static string SerializePayload(IServiceProvider serviceProvider, object obj) + { + var json = JsonSerializer.Serialize(obj, obj.GetType(), SerializerOptions); + return json; + } + + public static object DeserializePayload(IServiceProvider serviceProvider, string json, Type type) + { + var value = JsonSerializer.Deserialize(json, type, SerializerOptions)!; + return value; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Stimuli/MessageReceivedStimulus.cs b/src/modules/Elsa.Kafka/Stimuli/MessageReceivedStimulus.cs new file mode 100644 index 000000000..20c04921f --- /dev/null +++ b/src/modules/Elsa.Kafka/Stimuli/MessageReceivedStimulus.cs @@ -0,0 +1,14 @@ +using Elsa.Expressions.Models; +using Elsa.Workflows.Attributes; + +namespace Elsa.Kafka.Stimuli; + +public class MessageReceivedStimulus +{ + public string ConsumerDefinitionId { get; set; } = default!; + [ExcludeFromHash] public Type? MessageType { get; set; } + public ICollection Topics { get; set; } = []; + public IDictionary? CorrelatingFields { get; set; } + public Expression? Predicate { get; set; } + public bool IsLocal { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/Tasks/StartConsumersTask.cs b/src/modules/Elsa.Kafka/Tasks/StartConsumersTask.cs new file mode 100644 index 000000000..cc9d8c072 --- /dev/null +++ b/src/modules/Elsa.Kafka/Tasks/StartConsumersTask.cs @@ -0,0 +1,20 @@ +using Elsa.Common; +using JetBrains.Annotations; + +namespace Elsa.Kafka.Tasks; + +[UsedImplicitly] +public class StartConsumersStartupTask(IWorkerManager workerManager, IWorkerTopicSubscriber workerTopicSubscriber) : BackgroundTask +{ + public override async Task StartAsync(CancellationToken cancellationToken) + { + await workerManager.UpdateWorkersAsync(cancellationToken); + await workerTopicSubscriber.UpdateTopicSubscriptionsAsync(cancellationToken); + } + + public override Task StopAsync(CancellationToken cancellationToken) + { + workerManager.StopWorkers(); + return Task.CompletedTask; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/UIHints/ConsumerDefinitionsDropdownOptionsProvider.cs b/src/modules/Elsa.Kafka/UIHints/ConsumerDefinitionsDropdownOptionsProvider.cs new file mode 100644 index 000000000..ca1550f22 --- /dev/null +++ b/src/modules/Elsa.Kafka/UIHints/ConsumerDefinitionsDropdownOptionsProvider.cs @@ -0,0 +1,14 @@ +using System.Reflection; +using Elsa.Workflows.UIHints.Dropdown; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.UIHints; + +public class ConsumerDefinitionsDropdownOptionsProvider(IConsumerDefinitionEnumerator consumerEnumerator) : DropDownOptionsProviderBase +{ + protected override async ValueTask> GetItemsAsync(PropertyInfo propertyInfo, object? context, CancellationToken cancellationToken) + { + var definitions = await consumerEnumerator.EnumerateAsync(cancellationToken).ToList(); + return definitions.Select(x => new SelectListItem(x.Name, x.Id)).ToList(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/UIHints/ProducerDefinitionsDropdownOptionsProvider.cs b/src/modules/Elsa.Kafka/UIHints/ProducerDefinitionsDropdownOptionsProvider.cs new file mode 100644 index 000000000..af72114d7 --- /dev/null +++ b/src/modules/Elsa.Kafka/UIHints/ProducerDefinitionsDropdownOptionsProvider.cs @@ -0,0 +1,14 @@ +using System.Reflection; +using Elsa.Workflows.UIHints.Dropdown; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.UIHints; + +public class ProducerDefinitionsDropdownOptionsProvider(IProducerDefinitionEnumerator producerEnumerator) : DropDownOptionsProviderBase +{ + protected override async ValueTask> GetItemsAsync(PropertyInfo propertyInfo, object? context, CancellationToken cancellationToken) + { + var definitions = await producerEnumerator.EnumerateAsync(cancellationToken).ToList(); + return definitions.Select(x => new SelectListItem(x.Name, x.Id)).ToList(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Kafka/UIHints/TopicDefinitionsDropdownOptionsProvider.cs b/src/modules/Elsa.Kafka/UIHints/TopicDefinitionsDropdownOptionsProvider.cs new file mode 100644 index 000000000..bf029d912 --- /dev/null +++ b/src/modules/Elsa.Kafka/UIHints/TopicDefinitionsDropdownOptionsProvider.cs @@ -0,0 +1,14 @@ +using System.Reflection; +using Elsa.Workflows.UIHints.Dropdown; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Kafka.UIHints; + +public class TopicDefinitionsDropdownOptionsProvider(ITopicDefinitionEnumerator topicEnumerator) : DropDownOptionsProviderBase +{ + protected override async ValueTask> GetItemsAsync(PropertyInfo propertyInfo, object? context, CancellationToken cancellationToken) + { + var definitions = await topicEnumerator.EnumerateAsync(cancellationToken).ToList(); + return definitions.Select(x => new SelectListItem(x.Name, x.Id)).ToList(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Services/MassTransitActivityTypeProvider.cs b/src/modules/Elsa.MassTransit/Services/MassTransitActivityTypeProvider.cs index e8a8b0612..fde91ef3c 100644 --- a/src/modules/Elsa.MassTransit/Services/MassTransitActivityTypeProvider.cs +++ b/src/modules/Elsa.MassTransit/Services/MassTransitActivityTypeProvider.cs @@ -61,6 +61,7 @@ public class MassTransitActivityTypeProvider(IActivityFactory activityFactory, I return new() { + Name = typeName, TypeName = fullTypeName, Version = 1, DisplayName = displayName, diff --git a/src/modules/Elsa.Python/Features/PythonFeature.cs b/src/modules/Elsa.Python/Features/PythonFeature.cs index 79b5e39ad..32dbe5df1 100644 --- a/src/modules/Elsa.Python/Features/PythonFeature.cs +++ b/src/modules/Elsa.Python/Features/PythonFeature.cs @@ -4,11 +4,13 @@ using Elsa.Extensions; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; +using Elsa.Python.Activities; using Elsa.Python.Contracts; using Elsa.Python.HostedServices; using Elsa.Python.Options; using Elsa.Python.Providers; using Elsa.Python.Services; +using Elsa.Workflows; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Python.Features; @@ -52,5 +54,8 @@ public class PythonFeature : FeatureBase // Activities. Module.AddActivitiesFrom(); + + // UI property handlers. + Services.AddScoped(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs index 6123ce00e..bd8623ce8 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs @@ -533,6 +533,7 @@ public partial class WorkflowExecutionContext : IExecutionContext var parentContext = options?.Owner; var parentExpressionExecutionContext = parentContext?.ExpressionExecutionContext ?? ExpressionExecutionContext; var properties = ExpressionExecutionContextExtensions.CreateActivityExecutionContextPropertiesFrom(this, Input); + properties[ExpressionExecutionContextExtensions.ActivityKey] = activity; var memory = new MemoryRegister(); var now = SystemClock.UtcNow; var expressionExecutionContext = new ExpressionExecutionContext(ServiceProvider, memory, parentExpressionExecutionContext, properties, CancellationToken); diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IPayloadSerializer.cs b/src/modules/Elsa.Workflows.Core/Contracts/IPayloadSerializer.cs index 2fca52dc1..c2ccbaba7 100644 --- a/src/modules/Elsa.Workflows.Core/Contracts/IPayloadSerializer.cs +++ b/src/modules/Elsa.Workflows.Core/Contracts/IPayloadSerializer.cs @@ -28,6 +28,14 @@ public interface IPayloadSerializer /// The deserialized state. object Deserialize(string serializedData); + /// + /// Deserializes the specified serialized state. + /// + /// The serialized state. + /// The type to deserialize the state into. + /// The deserialized state. + object Deserialize(string serializedData, Type type); + /// /// Deserializes the specified serialized state. /// diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs index d796a1c18..8614afb6f 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ExpressionExecutionContextExtensions.cs @@ -37,6 +37,11 @@ public static class ExpressionExecutionContextExtensions /// The key used to store the workflow in the dictionary. /// public static readonly object WorkflowKey = new(); + + /// + /// The key used to store the activity in the dictionary. + /// + public static readonly object ActivityKey = new(); /// /// Creates a dictionary for the specified and . @@ -78,6 +83,11 @@ public static class ExpressionExecutionContextExtensions /// Returns the of the specified /// public static bool TryGetActivityExecutionContext(this ExpressionExecutionContext context, out ActivityExecutionContext activityExecutionContext) => context.TransientProperties.TryGetValue(ActivityExecutionContextKey, out activityExecutionContext!); + + /// + /// Returns the of the specified + /// + public static IActivity GetActivity(this ExpressionExecutionContext context) => (IActivity)context.TransientProperties[ActivityKey]; /// /// Returns the value of the specified input. diff --git a/src/modules/Elsa.Workflows.Core/Extensions/TriggerExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/TriggerExtensions.cs index 5b173e18e..a6de0a24f 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/TriggerExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/TriggerExtensions.cs @@ -2,6 +2,8 @@ using Elsa.Expressions.Contracts; using Elsa.Expressions.Models; using Elsa.Workflows; using Elsa.Workflows.Helpers; +using Elsa.Workflows.Models; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; // ReSharper disable once CheckNamespace @@ -22,39 +24,57 @@ public static class TriggerExtensions var activityTypeName = TypeNameHelper.GenerateTypeName(); return triggers.Where(x => x.Type == activityTypeName); } - + /// /// Creates an expression execution context for the specified trigger. /// /// The trigger for which to create an expression execution context. + /// The activity descriptor. /// The service provider. /// The workflow indexing context. /// The expression evaluator. /// The logger. /// An expression execution context. public static async Task CreateExpressionExecutionContextAsync( - this ITrigger trigger, - IServiceProvider serviceProvider, - WorkflowIndexingContext context, + this ITrigger trigger, + ActivityDescriptor activityDescriptor, + IServiceProvider serviceProvider, + WorkflowIndexingContext context, IExpressionEvaluator expressionEvaluator, ILogger logger) { - var inputs = trigger.GetInputs(); - var assignedInputs = inputs.Where(x => x.MemoryBlockReference != null!).ToList(); + var namedInputs = trigger.GetNamedInputs(); + var assignedInputs = namedInputs.Where(x => x.Value.MemoryBlockReference != null!).ToList(); var register = context.GetOrCreateRegister(trigger); var cancellationToken = context.CancellationToken; var expressionInput = new Dictionary(); var applicationProperties = ExpressionExecutionContextExtensions.CreateTriggerIndexingPropertiesFrom(context.Workflow, expressionInput); + applicationProperties[ExpressionExecutionContextExtensions.ActivityKey] = trigger; var expressionExecutionContext = new ExpressionExecutionContext(serviceProvider, register, default, applicationProperties, cancellationToken); // Evaluate activity inputs before requesting trigger data. - foreach (var input in assignedInputs) + foreach (var namedInput in assignedInputs) { + var inputDescriptor = activityDescriptor.Inputs.FirstOrDefault(x => x.Name == namedInput.Key); + + if (inputDescriptor == null) + { + logger.LogWarning("Input descriptor not found for input '{InputName}'", namedInput.Key); + continue; + } + + if (!inputDescriptor.AutoEvaluate) + { + logger.LogDebug("Skipping input '{InputName}' because it is not set to auto-evaluate.", namedInput.Key); + continue; + } + + var input = namedInput.Value; var locationReference = input.MemoryBlockReference(); - if(locationReference.Id == null!) + if (locationReference.Id == null!) continue; - + try { var value = await expressionEvaluator.EvaluateAsync(input, expressionExecutionContext); diff --git a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs index ada38c584..67e5cb4a6 100644 --- a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs +++ b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs @@ -1,6 +1,7 @@ using Elsa.Common; using Elsa.Common.Features; using Elsa.Common.Serialization; +using Elsa.CSharp.Activities; using Elsa.Expressions.Features; using Elsa.Extensions; using Elsa.Features.Abstractions; @@ -230,10 +231,15 @@ public class WorkflowsFeature : FeatureBase .AddScoped() .AddScoped() + // UI property handlers. + .AddScoped() + .AddScoped() + .AddScoped() + // Logger state generators. .AddSingleton(WorkflowLoggerStateGenerator) .AddSingleton(ActivityLoggerStateGenerator) - + // Log Persistence Strategies. .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Core/Models/Bookmark.cs b/src/modules/Elsa.Workflows.Core/Models/Bookmark.cs index d9337f5bf..9d062b336 100644 --- a/src/modules/Elsa.Workflows.Core/Models/Bookmark.cs +++ b/src/modules/Elsa.Workflows.Core/Models/Bookmark.cs @@ -23,7 +23,7 @@ public record Bookmark( object? Payload, string ActivityId, string ActivityNodeId, - string ActivityInstanceId, + string? ActivityInstanceId, DateTimeOffset CreatedAt, bool AutoBurn = true, string? CallbackMethodName = default, diff --git a/src/modules/Elsa.Workflows.Core/Serialization/Converters/ExpandoObjectConverter.cs b/src/modules/Elsa.Workflows.Core/Serialization/Converters/ExpandoObjectConverter.cs index c4c3229bc..b9b31cc55 100644 --- a/src/modules/Elsa.Workflows.Core/Serialization/Converters/ExpandoObjectConverter.cs +++ b/src/modules/Elsa.Workflows.Core/Serialization/Converters/ExpandoObjectConverter.cs @@ -12,7 +12,7 @@ public sealed class ExpandoObjectConverter : JsonConverter /// public override void Write(Utf8JsonWriter writer, object value, JsonSerializerOptions options) { - JsonSerializer.Serialize(writer, value, typeof(ExpandoObject), options); + JsonSerializer.Serialize(writer, value, value.GetType(), options); } /// diff --git a/src/modules/Elsa.Workflows.Core/Serialization/Serializers/JsonPayloadSerializer.cs b/src/modules/Elsa.Workflows.Core/Serialization/Serializers/JsonPayloadSerializer.cs index 90744a0ca..26b811d89 100644 --- a/src/modules/Elsa.Workflows.Core/Serialization/Serializers/JsonPayloadSerializer.cs +++ b/src/modules/Elsa.Workflows.Core/Serialization/Serializers/JsonPayloadSerializer.cs @@ -41,6 +41,12 @@ public class JsonPayloadSerializer : IPayloadSerializer return Deserialize(payload); } + public object Deserialize(string serializedData, Type type) + { + var options = GetOptions(); + return JsonSerializer.Deserialize(serializedData, type, options)!; + } + /// public object Deserialize(JsonElement payload) { diff --git a/src/modules/Elsa.Workflows.Core/Services/PropertyUIHandlerResolver.cs b/src/modules/Elsa.Workflows.Core/Services/PropertyUIHandlerResolver.cs index 236d925da..12feb3cab 100644 --- a/src/modules/Elsa.Workflows.Core/Services/PropertyUIHandlerResolver.cs +++ b/src/modules/Elsa.Workflows.Core/Services/PropertyUIHandlerResolver.cs @@ -2,48 +2,53 @@ using Elsa.Workflows.Attributes; using Microsoft.Extensions.DependencyInjection; -namespace Elsa.Workflows; - +namespace Elsa.Workflows; + /// -public class PropertyUIHandlerResolver(IServiceScopeFactory scopeFactory) : IPropertyUIHandlerResolver -{ +public class PropertyUIHandlerResolver(IServiceScopeFactory scopeFactory) : IPropertyUIHandlerResolver +{ /// - public async ValueTask> GetUIPropertiesAsync(PropertyInfo propertyInfo, object? context, CancellationToken cancellationToken = default) - { - var inputAttribute = propertyInfo.GetCustomAttribute(); - var result = new Dictionary(); - - var uiHandlers = new List(); - - if (inputAttribute?.UIHandler != null) - uiHandlers.Add(inputAttribute.UIHandler); - - if (inputAttribute?.UIHandlers != null) - uiHandlers.AddRange(inputAttribute.UIHandlers); - - using var scope = scopeFactory.CreateScope(); - var uiHintHandlers = scope.ServiceProvider.GetServices(); - - if (!string.IsNullOrWhiteSpace(inputAttribute?.UIHint)) - { - var uiHintHandler = uiHintHandlers.FirstOrDefault(x => x.UIHint == inputAttribute.UIHint); - - if (uiHintHandler != null) - { - var defaultHandlers = await uiHintHandler.GetPropertyUIHandlersAsync(propertyInfo, cancellationToken); - uiHandlers.AddRange(defaultHandlers); - } - } - - foreach (var handlerType in uiHandlers) - { - var provider = (IPropertyUIHandler)ActivatorUtilities.GetServiceOrCreateInstance(scope.ServiceProvider, handlerType); - var properties = await provider.GetUIPropertiesAsync(propertyInfo, context, cancellationToken); - - foreach (var property in properties) - result[property.Key] = property.Value; - } - - return result; - } + public async ValueTask> GetUIPropertiesAsync(PropertyInfo propertyInfo, object? context, CancellationToken cancellationToken = default) + { + var inputAttribute = propertyInfo.GetCustomAttribute(); + var result = new Dictionary(); + + var uiHandlers = new List(); + + if (inputAttribute?.UIHandler != null) + uiHandlers.Add(inputAttribute.UIHandler); + + if (inputAttribute?.UIHandlers != null) + uiHandlers.AddRange(inputAttribute.UIHandlers); + + using var scope = scopeFactory.CreateScope(); + var uiHintHandlers = scope.ServiceProvider.GetServices(); + + if (!string.IsNullOrWhiteSpace(inputAttribute?.UIHint)) + { + var uiHintHandler = uiHintHandlers.FirstOrDefault(x => x.UIHint == inputAttribute.UIHint); + + if (uiHintHandler != null) + { + var defaultHandlers = await uiHintHandler.GetPropertyUIHandlersAsync(propertyInfo, cancellationToken); + uiHandlers.AddRange(defaultHandlers); + } + } + + var propertyUIHandlers = scope.ServiceProvider.GetServices().ToList(); + foreach (var handlerType in uiHandlers) + { + var provider = propertyUIHandlers.FirstOrDefault(x => x.GetType() == handlerType); + + if (provider == null) + continue; + + var properties = await provider.GetUIPropertiesAsync(propertyInfo, context, cancellationToken); + + foreach (var property in properties) + result[property.Key] = property.Value; + } + + return result; + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Providers/DefaultExpressionDescriptorProvider.cs b/src/modules/Elsa.Workflows.Management/Providers/DefaultExpressionDescriptorProvider.cs index 67d1994f8..a07a0e393 100644 --- a/src/modules/Elsa.Workflows.Management/Providers/DefaultExpressionDescriptorProvider.cs +++ b/src/modules/Elsa.Workflows.Management/Providers/DefaultExpressionDescriptorProvider.cs @@ -53,19 +53,21 @@ public class DefaultExpressionDescriptorProvider : IExpressionDescriptorProvider memoryBlockReferenceFactory: () => new Variable(), deserialize: context => { - var valueElement = context.JsonElement.TryGetProperty("value", out var v) ? v : default; - var valueString = valueElement.GetValue()?.ToString(); - try - { - var value = JsonSerializer.Deserialize(valueString, context.MemoryBlockType, context.Options); - return new Expression("Variable", value); - } - catch (Exception) - { - - return new Expression("Variable", null); - } - + var valueElement = context.JsonElement.TryGetProperty("value", out var v) ? v : default; + var valueString = valueElement.GetValue()?.ToString(); + + if (string.IsNullOrWhiteSpace(valueString)) + return new Expression("Variable", null); + + try + { + var value = JsonSerializer.Deserialize(valueString, context.MemoryBlockType, context.Options); + return new Expression("Variable", value); + } + catch (Exception) + { + return new Expression("Variable", null); + } } ); } diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs index c4484d6ea..865553096 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs @@ -27,6 +27,9 @@ public interface IBookmarkResumer /// Resumes the bookmark. /// Task ResumeAsync(string bookmarkId, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity; + + /// Resumes the bookmark. + Task ResumeAsync(ResumeBookmarkRequest request, CancellationToken cancellationToken = default); /// /// Resumes the bookmark. diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerInvoker.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerInvoker.cs new file mode 100644 index 000000000..5692d8986 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerInvoker.cs @@ -0,0 +1,6 @@ +namespace Elsa.Workflows.Runtime; + +public interface ITriggerInvoker +{ + Task InvokeAsync(InvokeTriggerRequest request, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/StoredBookmark.cs b/src/modules/Elsa.Workflows.Runtime/Entities/StoredBookmark.cs index 0b11690f9..e92411220 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/StoredBookmark.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/StoredBookmark.cs @@ -1,4 +1,5 @@ using Elsa.Common.Entities; +using Elsa.Workflows.Models; namespace Elsa.Workflows.Runtime.Entities; @@ -46,4 +47,18 @@ public class StoredBookmark : Entity /// The date and time the bookmark was created. /// public DateTimeOffset CreatedAt { get; set; } + + public Bookmark ToBookmark() + { + return new Bookmark + { + Id = Id, + Name = ActivityTypeName, + Hash = Hash, + CreatedAt = CreatedAt, + ActivityInstanceId = ActivityInstanceId, + Payload = Payload, + Metadata = Metadata + }; + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index 5b604a6f1..5adcb1735 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -20,6 +20,7 @@ using Elsa.Workflows.Runtime.Providers; using Elsa.Workflows.Runtime.Services; using Elsa.Workflows.Runtime.Stores; using Elsa.Workflows.Runtime.Tasks; +using Elsa.Workflows.Runtime.UIHints; using Medallion.Threading; using Medallion.Threading.FileSystem; using Microsoft.Extensions.DependencyInjection; @@ -251,6 +252,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() @@ -297,6 +299,9 @@ public class WorkflowRuntimeFeature : FeatureBase // Workflow definition providers. .AddWorkflowDefinitionProvider() + + // UI prooprty handlers. + .AddScoped() // Domain handlers. .AddCommandHandler() diff --git a/src/modules/Elsa.Workflows.Runtime/Requests/InvokeTriggerRequest.cs b/src/modules/Elsa.Workflows.Runtime/Requests/InvokeTriggerRequest.cs new file mode 100644 index 000000000..c2580b187 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Requests/InvokeTriggerRequest.cs @@ -0,0 +1,14 @@ +using Elsa.Workflows.Activities; + +namespace Elsa.Workflows.Runtime; + +public class InvokeTriggerRequest +{ + public Workflow Workflow { get; set; } = default!; + public string ActivityId { get; set; } = default!; + public string? CorrelationId { get; set; } + public IDictionary? Input { get; set; } + public IDictionary? Properties { get; set; } + public string? ParentWorkflowInstanceId { get; set; } + +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs b/src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs new file mode 100644 index 000000000..58a29bed6 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs @@ -0,0 +1,20 @@ +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Runtime; + +public class ResumeBookmarkRequest +{ + public string WorkflowInstanceId { get; set; } = default!; + + /// The ID of the bookmark that triggered the workflow instance, if any. + public string BookmarkId { get; set; } = default!; + + /// The handle of the activity to schedule, if any. + public ActivityHandle? ActivityHandle { get; set; } + + /// Any additional properties to associate with the workflow instance. + public IDictionary? Properties { get; set; } + + /// The input to the workflow instance, if any. + public IDictionary? Input { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs index 519b9ef1c..fa0ead919 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs @@ -55,6 +55,22 @@ public class BookmarkResumer(IWorkflowRuntime workflowRuntime, IBookmarkStore bo return await ResumeAsync(bookmarkFilter, options, cancellationToken); } + public async Task ResumeAsync(ResumeBookmarkRequest request, CancellationToken cancellationToken = default) + { + var runRequest = new RunWorkflowInstanceRequest + { + Input = request.Input, + Properties = request.Properties, + ActivityHandle = request.ActivityHandle, + BookmarkId = request.BookmarkId + }; + + var workflowInstanceId = request.WorkflowInstanceId; + var workflowClient = await workflowRuntime.CreateClientAsync(workflowInstanceId, cancellationToken); + var response = await workflowClient.RunInstanceAsync(runRequest, cancellationToken); + return ResumeBookmarkResult.Found(response); + } + /// public async Task ResumeAsync(BookmarkFilter filter, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs b/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs index 3d4a717ee..7ded4d3a3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs @@ -15,6 +15,7 @@ public class StimulusSender( IBookmarkQueue bookmarkQueue, IWorkflowRuntime workflowRuntime, IWorkflowStarter workflowStarter, + ITriggerInvoker triggerInvoker, ILogger logger) : IStimulusSender { /// @@ -56,24 +57,25 @@ public class StimulusSender( foreach (var trigger in triggerBoundWorkflow.Triggers) { - var startRequest = new StartWorkflowRequest + var triggerRequest = new InvokeTriggerRequest { - Workflow = workflow, CorrelationId = correlationId, + Workflow = workflow, + ActivityId = trigger.ActivityId, Input = input, Properties = properties, - ParentId = parentId, - TriggerActivityId = trigger.ActivityId + ParentWorkflowInstanceId = parentId }; - var startResponse = await workflowStarter.StartWorkflowAsync(startRequest, cancellationToken); - if (startResponse.CannotStart) + var response = await triggerInvoker.InvokeAsync(triggerRequest, cancellationToken); + + if (response.CannotStart) { logger.LogWarning("Workflow activation strategy disallowed starting workflow {WorkflowDefinitionHandle} with correlation ID {CorrelationId}", workflow.DefinitionHandle, correlationId); continue; } - responses.Add(startResponse.ToRunWorkflowInstanceResponse()); + responses.Add(response.ToRunWorkflowInstanceResponse()); } } @@ -119,27 +121,21 @@ public class StimulusSender( } else { - // If no bookmarks were matched, but bookmark-specific details were given, enqueue the request in case a matching bookmark is created in the near future. + // If no bookmarks were matched, enqueue the request in case a matching bookmark is created in the near future. var workflowInstanceId = metadata?.WorkflowInstanceId; - var bookmarkId = metadata?.BookmarkId; - var activityInstanceId = metadata?.ActivityInstanceId; - var correlationId = metadata?.CorrelationId; - if (workflowInstanceId != null || bookmarkId != null || activityInstanceId != null || correlationId != null) + var bookmarkQueueItem = new NewBookmarkQueueItem { - var bookmarkQueueItem = new NewBookmarkQueueItem + WorkflowInstanceId = workflowInstanceId, + BookmarkId = metadata?.BookmarkId, + StimulusHash = stimulusHash, + Options = new ResumeBookmarkOptions { - WorkflowInstanceId = workflowInstanceId, - BookmarkId = metadata?.BookmarkId, - StimulusHash = stimulusHash, - Options = new ResumeBookmarkOptions - { - Input = input, - Properties = properties - } - }; - await bookmarkQueue.EnqueueAsync(bookmarkQueueItem, cancellationToken); - } + Input = input, + Properties = properties + } + }; + await bookmarkQueue.EnqueueAsync(bookmarkQueueItem, cancellationToken); } return responses; diff --git a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs index a9ee989d8..4e3d4cc96 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs @@ -23,6 +23,7 @@ public class TriggerIndexer : ITriggerIndexer private readonly IExpressionEvaluator _expressionEvaluator; private readonly IIdentityGenerator _identityGenerator; private readonly ITriggerStore _triggerStore; + private readonly IActivityRegistry _activityRegistry; private readonly INotificationSender _notificationSender; private readonly IServiceProvider _serviceProvider; private readonly IStimulusHasher _hasher; @@ -37,6 +38,7 @@ public class TriggerIndexer : ITriggerIndexer IExpressionEvaluator expressionEvaluator, IIdentityGenerator identityGenerator, ITriggerStore triggerStore, + IActivityRegistry activityRegistry, INotificationSender notificationSender, IServiceProvider serviceProvider, IStimulusHasher hasher, @@ -46,6 +48,7 @@ public class TriggerIndexer : ITriggerIndexer _expressionEvaluator = expressionEvaluator; _identityGenerator = identityGenerator; _triggerStore = triggerStore; + _activityRegistry = activityRegistry; _notificationSender = notificationSender; _serviceProvider = serviceProvider; _hasher = hasher; @@ -153,11 +156,18 @@ public class TriggerIndexer : ITriggerIndexer { var workflow = context.Workflow; var cancellationToken = context.CancellationToken; - var expressionExecutionContext = await trigger.CreateExpressionExecutionContextAsync(_serviceProvider, context, _expressionEvaluator, _logger); - + var triggerTypeName = trigger.Type; + var triggerDescriptor = _activityRegistry.Find(triggerTypeName, trigger.Version); + + if (triggerDescriptor == null) + { + _logger.LogWarning("Could not find activity descriptor for activity type {ActivityType}", triggerTypeName); + return new List(0); + } + + var expressionExecutionContext = await trigger.CreateExpressionExecutionContextAsync(triggerDescriptor, _serviceProvider, context, _expressionEvaluator, _logger); var triggerIndexingContext = new TriggerIndexingContext(context, expressionExecutionContext, trigger, cancellationToken); var triggerData = await TryGetTriggerDataAsync(trigger, triggerIndexingContext); - var triggerTypeName = trigger.Type; // If no trigger payloads were returned, create a null payload. if (!triggerData.Any()) triggerData.Add(null!); diff --git a/src/modules/Elsa.Workflows.Runtime/Services/TriggerInvoker.cs b/src/modules/Elsa.Workflows.Runtime/Services/TriggerInvoker.cs new file mode 100644 index 000000000..8d5bb2d9d --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/TriggerInvoker.cs @@ -0,0 +1,19 @@ +namespace Elsa.Workflows.Runtime; + +public class TriggerInvoker(IWorkflowStarter workflowStarter) : ITriggerInvoker +{ + public async Task InvokeAsync(InvokeTriggerRequest request, CancellationToken cancellationToken = default) + { + var startRequest = new StartWorkflowRequest + { + Workflow = request.Workflow, + CorrelationId = request.CorrelationId, + Input = request.Input, + Properties = request.Properties, + ParentId = request.ParentWorkflowInstanceId, + TriggerActivityId = request.ActivityId, + }; + + return await workflowStarter.StartWorkflowAsync(startRequest, cancellationToken); + } +} \ No newline at end of file