From f53e024d2565647d99cd4b0d4ced13ebebcc1ebf Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 18 Nov 2024 13:42:54 +0100 Subject: [PATCH] Add Elsa.Kafka Module for Kafka Integration with Message Sending and Receiving Activities (#6108) * Add Kafka module with integration and example setup Introduced the Kafka module providing consumer integration and activities into the project. This includes new classes for consumer handling, configuration, and activities. An example setup using Docker Compose is also added to facilitate development and testing. * Refactor KafkaTransportMessage to inline Timestamp namespace Simplify the namespace usage for the Timestamp type within the KafkaTransportMessage record. This change eliminates the need for a separate using directive for Timestamp, enhancing code readability and maintainability. * Add support for handling Kafka transport messages This commit introduces the capability to handle and trigger workflows based on Kafka transport messages. It adds a new handler, notifications, and updates the message stimulus to include correlating fields. Additionally, the Kafka consumers are now managed more modularly with updated startup tasks and mediator integration. * Enable Kafka integration and fix Kafka options naming Added support for Kafka integration in Elsa.Server.Web by setting up Kafka configurations in appsettings.json and updating Program.cs. Also, renamed `ConsumerConfigs` to `ConsumerDefinitions` in Kafka options for clarity. * Add consumer definition enumeration and dropdown support Introduced `IConsumerDefinitionEnumerator` for managing consumer definitions across providers and implemented in `ConsumerDefinitionEnumerator` class. Enhanced `KafkaFeature` to register these services and updated the `MessageReceived` activity to use a dropdown UI hint for consuming definitions. Improved `StartConsumersTask` by refactoring consumer definition retrieval logic. * Add SendMessage activity and refine Kafka messaging Introduce a new SendMessage activity for Kafka, enabling message publishing to specific topics. Refine KafkaTransportMessage model by removing headers and timestamp fields. Adjust the StimulusSender logic to streamline the bookmark queuing process and fix key-value pairing in dropdown options. Update appsettings for corrected Kafka bootstrap server and topic configurations. * Add producer and topic management support Introduced interfaces and implementations for managing producer and topic definitions along with their respective enumerators and list providers. Updated `SendMessage` activity to include producer selection and refactored consumer definition providers for better consistency. * Refactor Kafka configuration property names Renamed Kafka configuration properties for better consistency and readability across the codebase. Updated property names from `ProducerDefinitions` to `Producers`, `ConsumerDefinitions` to `Consumers`, and `TopicDefinitions` to `Topics`. Added missing input attribute in `SendMessage.cs` and registered additional handlers in `KafkaFeature.cs`. * Add custom serializers for Kafka message handling Introduced `DefaultSerializers` class for custom serialization and deserialization of Kafka messages. Updated `MessageReceived` and `SendMessage` activities to use these custom serializers, and modified `KafkaOptions` to include them. * Fix ExpandoObject serialization method parameter Changed the serialization type from `ExpandoObject` to the actual type of the object to ensure proper serialization. This ensures that derived types are correctly handled during the serialization process. * Add Producer and Consumer workflows for Kafka Introduced two new workflows: `ProducerWorkflow` and `ConsumerWorkflow` for handling Kafka messages. Updated `DefaultSerializers` to use camelCase property naming and modified `appsettings.json` to include `topic-2` and format entries. * Add JSON serialization to log output in ConsumerWorkflow This change enhances the log output by serializing messages to JSON format before writing them. The addition of System.Text.Json ensures that the message content is presented in a structured and standardized format in logs. * Add correlation strategies and update Kafka features Implemented HeaderCorrelationStrategy and NullCorrelationStrategy, and updated KafkaFeature to support customizable correlation strategies. Added correlation ID handling to Kafka transport messages and updated config and handlers accordingly. * Add tenant accessor to ConsumerDefinitionWorkflowContextProvider Integrated ITenantAccessor to the provider to support tenant-specific context loading. Updated the constructor and LoadAsync method to retrieve the tenant information and use it for context-specific operations. * Switch to MySQL and disable Kafka This commit changes the SQL database provider from SQLite to MySQL and disables Kafka use. It also includes necessary adjustments such as adding MySQL handling in configuration and connection setups, updating `docker-compose` to include MySQL services, and referencing MySQL projects in the `.csproj` file. * Add UI property handlers to multiple features This commit introduces various UI property handlers across several features such as Python, JavaScript, CSharp, and Workflow features to enhance user interface property handling. It also updates the property UI handler resolution logic to better manage cases where providers are not available. Furthermore, adjustments were made in the server configuration to switch database providers and enable Kafka. * Refactor property UI handler retrieval logic Modified the logic to fetch property UI handlers by preloading them into a list and then filtering. This change improves readability and potentially performance by reducing repetitive service provider calls. * Incremental work on Kafka workers and predicate evaluation * Merge BookmarkInvoker with BookmarkResumer * Register IWorkerManager * Change lifetime scope of WorkerManager to Singleton * **Introduce topic subscription handling for Kafka workers** Added `IWorkerTopicSubscriber` interface and its implementation for managing topic subscriptions. Enhanced workers to bind triggers and bookmarks dynamically based on existing data. Updated several classes and methods to support topic-based subscriptions and headers. * Refactor trigger matching logic. Extract trigger matching conditions into `IsMatchAsync` method for reuse. This enhances code maintainability and readability by reducing redundancy. The new `GetTopic` helper method isolates the topic retrieval logic. * Add Name property to MassTransitActivityTypeProvider This commit inserts the Name property in the returned object within the MassTransitActivityTypeProvider class. It ensures that the typeName is included, providing a clearer definition of the activity type. * Add handling for deleted bookmarks and refactor bookmark removal Added a new event handler for `BookmarksDeleted` to ensure removed bookmarks are processed correctly. Refactored the bookmark removal logic into a helper method to reduce code duplication and streamline the workflow. * Switch to asynchronous bookmark queue processing Refactored the `TriggerWorkflows` handler to use `IBookmarkQueue` instead of directly invoking the `IBookmarkResumer`. This change aims to improve scalability by queueing bookmark resumption requests, enabling better load distribution and async processing. Added necessary helpers and configuration options to support this functionality. * Remove unused IBookmarkResumer dependency Simplify the constructor by removing the unused IBookmarkResumer dependency. This cleanup reduces potential confusion and improves code maintainability without impacting functionality. * Add support for local message processing Introduced an `IsLocal` property to `MessageReceivedStimulus` for determining if the message event is local to a specific workflow instance. Updated `BookmarkBinding` and related handler methods to utilize `CorrelationId` for local event matching. Removed unused `CorrelatingFields` from `MessageReceived` activity. * Add nullability checks to IWorker retrieval methods Updated `GetWorker` methods to return nullable `IWorker` to handle cases where a worker might not exist. Modified code to include null checks and conditional operations to prevent potential null reference exceptions when accessing worker methods. * Add filtering based on activity type name for triggers and bookmarks This commit introduces filtering for triggers and bookmarks based on the `MessageReceived` activity type name. It also adds an option to mark messages as local in the `SendMessage` activity, where local messages are delivered to the current workflow instance only. These changes help enhance the management and targeted delivery of messages within the workflow framework. * Add new Kafka topics and clean up producers config New topics "topic-3" and "topic-4" were added to the Kafka settings. Unused topic references were removed from the producers configuration to simplify and improve clarity. * Add predicate to KafkaConsumerActivity in ConsumerWorkflow Introduced a predicate to the KafkaConsumerActivity using JavaScript expressions to filter messages based on OrderId. This ensures only relevant messages are processed in the workflow. --- Directory.Packages.props | 3 +- Elsa.sln | 8 + docker/docker-compose-kafka.yml | 167 ++++++++++++++ docker/docker-compose.yml | 14 ++ .../Elsa.Server.Web/Elsa.Server.Web.csproj | 4 + src/apps/Elsa.Server.Web/Program.cs | 22 ++ .../ConsumerConfigWorkflowContextProvider.cs | 22 ++ .../Workflows/ConsumerWorkflow.cs | 33 +++ .../Workflows/ProducerWorkflow.cs | 22 ++ src/apps/Elsa.Server.Web/appsettings.json | 56 +++++ .../Elsa.CSharp/Features/CSharpFeature.cs | 5 + .../Models/ExpressionExecutionContext.cs | 34 +-- src/modules/Elsa.Http/Features/HttpFeature.cs | 1 + .../Features/JavaScriptFeature.cs | 10 +- .../Elsa.Kafka/Activities/MessageReceived.cs | 130 +++++++++++ .../Elsa.Kafka/Activities/SendMessage.cs | 102 +++++++++ .../IConsumerDefinitionEnumerator.cs | 11 + .../Contracts/IConsumerDefinitionProvider.cs | 6 + .../Contracts/ICorrelationStrategy.cs | 9 + .../IProducerDefinitionEnumerator.cs | 16 ++ .../Contracts/IProducerDefinitionProvider.cs | 6 + .../Contracts/ITopicDefinitionEnumerator.cs | 11 + .../Contracts/ITopicDefinitionProvider.cs | 6 + src/modules/Elsa.Kafka/Contracts/IWorker.cs | 11 + .../Elsa.Kafka/Contracts/IWorkerManager.cs | 12 + .../Contracts/IWorkerTopicSubscriber.cs | 9 + .../Correlation/HeaderCorrelationStrategy.cs | 19 ++ .../Correlation/NullCorrelationStrategy.cs | 11 + src/modules/Elsa.Kafka/Elsa.Kafka.csproj | 24 ++ .../Elsa.Kafka/Elsa.Kafka.csproj.DotSettings | 9 + .../Entities/ConsumerConfigDefinition.cs | 13 ++ .../Elsa.Kafka/Entities/ProducerDefinition.cs | 9 + .../Elsa.Kafka/Entities/TopicDefinition.cs | 11 + .../Elsa.Kafka/Extensions/ModuleExtensions.cs | 12 + .../Elsa.Kafka/Features/KafkaFeature.cs | 76 +++++++ src/modules/Elsa.Kafka/FodyWeavers.xml | 3 + .../Elsa.Kafka/Handlers/TriggerWorkflows.cs | 208 ++++++++++++++++++ .../Elsa.Kafka/Handlers/UpdateWorkers.cs | 102 +++++++++ .../ConsumerDefinitionEnumerator.cs | 23 ++ .../ProducerDefinitionEnumerator.cs | 29 +++ .../TopicDefinitionEnumerator.cs | 23 ++ .../Elsa.Kafka/Implementations/Worker.cs | 128 +++++++++++ .../Implementations/WorkerManager.cs | 136 ++++++++++++ .../Implementations/WorkerTopicSubscriber.cs | 38 ++++ .../Elsa.Kafka/Models/BookmarkBinding.cs | 5 + .../Models/KafkaTransportMessage.cs | 3 + .../Elsa.Kafka/Models/TriggerBinding.cs | 6 + .../Notifications/TransportMessageReceived.cs | 6 + .../Elsa.Kafka/Options/KafkaOptions.cs | 16 ++ .../Providers/OptionsDefinitionProvider.cs | 21 ++ .../Serialization/DefaultSerializers.cs | 35 +++ .../Stimuli/MessageReceivedStimulus.cs | 14 ++ .../Elsa.Kafka/Tasks/StartConsumersTask.cs | 20 ++ ...sumerDefinitionsDropdownOptionsProvider.cs | 14 ++ ...ducerDefinitionsDropdownOptionsProvider.cs | 14 ++ ...TopicDefinitionsDropdownOptionsProvider.cs | 14 ++ .../MassTransitActivityTypeProvider.cs | 1 + .../Elsa.Python/Features/PythonFeature.cs | 5 + .../Contexts/WorkflowExecutionContext.cs | 1 + .../Contracts/IPayloadSerializer.cs | 8 + .../ExpressionExecutionContextExtensions.cs | 10 + .../Extensions/TriggerExtensions.cs | 38 +++- .../Features/WorkflowsFeature.cs | 8 +- .../Elsa.Workflows.Core/Models/Bookmark.cs | 2 +- .../Converters/ExpandoObjectConverter.cs | 2 +- .../Serializers/JsonPayloadSerializer.cs | 6 + .../Services/PropertyUIHandlerResolver.cs | 89 ++++---- .../DefaultExpressionDescriptorProvider.cs | 28 +-- .../Contracts/IBookmarkResumer.cs | 3 + .../Contracts/ITriggerInvoker.cs | 6 + .../Entities/StoredBookmark.cs | 15 ++ .../Features/WorkflowRuntimeFeature.cs | 5 + .../Requests/InvokeTriggerRequest.cs | 14 ++ .../Requests/ResumeBookmarkRequest.cs | 20 ++ .../Services/BookmarkResumer.cs | 16 ++ .../Services/StimulusSender.cs | 44 ++-- .../Services/TriggerIndexer.cs | 16 +- .../Services/TriggerInvoker.cs | 19 ++ 78 files changed, 2006 insertions(+), 122 deletions(-) create mode 100644 docker/docker-compose-kafka.yml create mode 100644 src/apps/Elsa.Server.Web/WorkflowContextProviders/ConsumerConfigWorkflowContextProvider.cs create mode 100644 src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs create mode 100644 src/apps/Elsa.Server.Web/Workflows/ProducerWorkflow.cs create mode 100644 src/modules/Elsa.Kafka/Activities/MessageReceived.cs create mode 100644 src/modules/Elsa.Kafka/Activities/SendMessage.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionEnumerator.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/IConsumerDefinitionProvider.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/ICorrelationStrategy.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/IProducerDefinitionEnumerator.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/IProducerDefinitionProvider.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/ITopicDefinitionEnumerator.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/ITopicDefinitionProvider.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/IWorker.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/IWorkerManager.cs create mode 100644 src/modules/Elsa.Kafka/Contracts/IWorkerTopicSubscriber.cs create mode 100644 src/modules/Elsa.Kafka/Correlation/HeaderCorrelationStrategy.cs create mode 100644 src/modules/Elsa.Kafka/Correlation/NullCorrelationStrategy.cs create mode 100644 src/modules/Elsa.Kafka/Elsa.Kafka.csproj create mode 100644 src/modules/Elsa.Kafka/Elsa.Kafka.csproj.DotSettings create mode 100644 src/modules/Elsa.Kafka/Entities/ConsumerConfigDefinition.cs create mode 100644 src/modules/Elsa.Kafka/Entities/ProducerDefinition.cs create mode 100644 src/modules/Elsa.Kafka/Entities/TopicDefinition.cs create mode 100644 src/modules/Elsa.Kafka/Extensions/ModuleExtensions.cs create mode 100644 src/modules/Elsa.Kafka/Features/KafkaFeature.cs create mode 100644 src/modules/Elsa.Kafka/FodyWeavers.xml create mode 100644 src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs create mode 100644 src/modules/Elsa.Kafka/Handlers/UpdateWorkers.cs create mode 100644 src/modules/Elsa.Kafka/Implementations/ConsumerDefinitionEnumerator.cs create mode 100644 src/modules/Elsa.Kafka/Implementations/ProducerDefinitionEnumerator.cs create mode 100644 src/modules/Elsa.Kafka/Implementations/TopicDefinitionEnumerator.cs create mode 100644 src/modules/Elsa.Kafka/Implementations/Worker.cs create mode 100644 src/modules/Elsa.Kafka/Implementations/WorkerManager.cs create mode 100644 src/modules/Elsa.Kafka/Implementations/WorkerTopicSubscriber.cs create mode 100644 src/modules/Elsa.Kafka/Models/BookmarkBinding.cs create mode 100644 src/modules/Elsa.Kafka/Models/KafkaTransportMessage.cs create mode 100644 src/modules/Elsa.Kafka/Models/TriggerBinding.cs create mode 100644 src/modules/Elsa.Kafka/Notifications/TransportMessageReceived.cs create mode 100644 src/modules/Elsa.Kafka/Options/KafkaOptions.cs create mode 100644 src/modules/Elsa.Kafka/Providers/OptionsDefinitionProvider.cs create mode 100644 src/modules/Elsa.Kafka/Serialization/DefaultSerializers.cs create mode 100644 src/modules/Elsa.Kafka/Stimuli/MessageReceivedStimulus.cs create mode 100644 src/modules/Elsa.Kafka/Tasks/StartConsumersTask.cs create mode 100644 src/modules/Elsa.Kafka/UIHints/ConsumerDefinitionsDropdownOptionsProvider.cs create mode 100644 src/modules/Elsa.Kafka/UIHints/ProducerDefinitionsDropdownOptionsProvider.cs create mode 100644 src/modules/Elsa.Kafka/UIHints/TopicDefinitionsDropdownOptionsProvider.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerInvoker.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Requests/InvokeTriggerRequest.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Services/TriggerInvoker.cs 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