From 44bc1d81af45b75951cee06be880b0730686ea8f Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 22 Nov 2024 20:47:44 +0100 Subject: [PATCH] Fix error handling and logging in Kafka components Enhanced error handling in Worker.cs to catch consume exceptions and log warnings. Added logging to TriggerWorkflows.cs for better traceability and fixed issues with expression evaluation. Also updated predicate in ConsumerWorkflow.cs and set CanStartWorkflow to true. --- .../Workflows/ConsumerWorkflow.cs | 4 +-- .../Elsa.Kafka/Handlers/TriggerWorkflows.cs | 27 ++++++++++++---- .../Elsa.Kafka/Implementations/Worker.cs | 32 +++++++++++++++---- 3 files changed, 47 insertions(+), 16 deletions(-) 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();