diff --git a/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs b/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs index 4d490ad7f..cf24baeb3 100644 --- a/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs +++ b/src/apps/Elsa.Server.Web/Workflows/ConsumerWorkflow.cs @@ -21,9 +21,9 @@ public class ConsumerWorkflow : WorkflowBase { ConsumerDefinitionId = new("consumer-1"), Topics = new(["topic-1"]), - Predicate = new(JavaScriptExpression.Create("message => message.OrderId == '1'")), + Predicate = new(JavaScriptExpression.Create("getMessage().OrderId == '1'")), Result = new(message), - CanStartWorkflow = false + CanStartWorkflow = true }, new WriteLine(c => JsonSerializer.Serialize(message.Get(c))) } diff --git a/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs b/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs index 41055fbfa..c11e45d7f 100644 --- a/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs +++ b/src/modules/Elsa.Kafka/Handlers/TriggerWorkflows.cs @@ -10,6 +10,7 @@ using Elsa.Workflows.Memory; using Elsa.Workflows.Runtime; using Elsa.Workflows.Runtime.Options; using JetBrains.Annotations; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; namespace Elsa.Kafka.Handlers; @@ -21,7 +22,8 @@ public class TriggerWorkflows( ICorrelationStrategy correlationStrategy, IExpressionEvaluator expressionEvaluator, IOptions options, - IServiceProvider serviceProvider) : INotificationHandler + IServiceProvider serviceProvider, + ILogger logger) : INotificationHandler { private static readonly string MessageReceivedActivityTypeName = ActivityTypeNameHelper.GenerateTypeName(); @@ -137,7 +139,7 @@ public class TriggerWorkflows( var correlationId = GetCorrelationId(transportMessage); var workflowInstanceId = GetWorkflowInstanceId(transportMessage); - + foreach (var binding in boundBookmarks) { var stimulus = binding.Stimulus; @@ -176,13 +178,24 @@ public class TriggerWorkflows( return true; var memory = new MemoryRegister(); - var messageVariable = new Variable("message", transportMessage); - var message = transportMessage; + var transportMessageVariable = new Variable("transportMessage", transportMessage); + var messageVariable = new Variable("message", transportMessage.Value); var expressionExecutionContext = new ExpressionExecutionContext(serviceProvider, memory, cancellationToken: cancellationToken); - messageVariable.Set(expressionExecutionContext, message); - return await expressionEvaluator.EvaluateAsync(predicate, expressionExecutionContext); + + transportMessageVariable.Set(expressionExecutionContext, transportMessage); + messageVariable.Set(expressionExecutionContext, transportMessage.Value); + + try + { + return await expressionEvaluator.EvaluateAsync(predicate, expressionExecutionContext); + } + catch (Exception e) + { + logger.LogWarning(e, "An error occurred while evaluating the predicate for stimulus {Stimulus}", stimulus); + return false; + } } - + private string? GetWorkflowInstanceId(KafkaTransportMessage transportMessage) { var key = options.Value.WorkflowInstanceIdHeaderKey; diff --git a/src/modules/Elsa.Kafka/Implementations/Worker.cs b/src/modules/Elsa.Kafka/Implementations/Worker.cs index 1e58d4ae3..218cdf23a 100644 --- a/src/modules/Elsa.Kafka/Implementations/Worker.cs +++ b/src/modules/Elsa.Kafka/Implementations/Worker.cs @@ -85,26 +85,44 @@ public class Worker(WorkerContext workerContext, IConsumer 100) + throw new InvalidOperationException("Too many consume exceptions."); + } + catch (OperationCanceledException) + { + logger.LogInformation("Consumer was cancelled."); + break; + } } consumer.Unsubscribe();