diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs index 226bb40e9..cc0a459be 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs @@ -13,12 +13,13 @@ using Elsa.Common.Entities; using Elsa.Workflows; using Elsa.Workflows.Management; using JetBrains.Annotations; +using Microsoft.Extensions.Logging; namespace Elsa.EntityFrameworkCore.Modules.Management; /// [UsedImplicitly] -public class EFCoreWorkflowDefinitionStore(EntityStore store, IPayloadSerializer payloadSerializer) +public class EFCoreWorkflowDefinitionStore(EntityStore store, IPayloadSerializer payloadSerializer, ILogger logger) : IWorkflowDefinitionStore { /// @@ -178,8 +179,15 @@ public class EFCoreWorkflowDefinitionStore(EntityStore(json); + try + { + if (!string.IsNullOrWhiteSpace(json)) + data = payloadSerializer.Deserialize(json); + } + catch (Exception exp) + { + logger.LogError(exp, "Could not deserialize workflow definition state: {DefinitionId}. Reverting to default state", entity.DefinitionId); + } entity.Options = data.Options; entity.Variables = data.Variables; diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs index 9aebbfaa4..79bee6bb7 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs @@ -11,6 +11,7 @@ using Elsa.Workflows.Management.Models; using Elsa.Workflows.Management.Options; using JetBrains.Annotations; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Open.Linq.AsyncExtensions; @@ -26,6 +27,7 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore private readonly IWorkflowStateSerializer _workflowStateSerializer; private readonly ICompressionCodecResolver _compressionCodecResolver; private readonly IOptions _options; + private readonly ILogger _logger; /// /// Constructor. @@ -34,12 +36,14 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore EntityStore store, IWorkflowStateSerializer workflowStateSerializer, ICompressionCodecResolver compressionCodecResolver, - IOptions options) + IOptions options, + ILogger logger) { _store = store; _workflowStateSerializer = workflowStateSerializer; _compressionCodecResolver = compressionCodecResolver; _options = options; + _logger = logger; } /// @@ -236,12 +240,18 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore var compressionAlgorithm = (string?)managementElsaDbContext.Entry(entity).Property("DataCompressionAlgorithm").CurrentValue ?? nameof(None); var compressionStrategy = _compressionCodecResolver.Resolve(compressionAlgorithm); - if (!string.IsNullOrWhiteSpace(json)) + try { - json = await compressionStrategy.DecompressAsync(json, cancellationToken); - data = _workflowStateSerializer.Deserialize(json); + if (!string.IsNullOrWhiteSpace(json)) + { + json = await compressionStrategy.DecompressAsync(json, cancellationToken); + data = _workflowStateSerializer.Deserialize(json); + } + } + catch (Exception exp) + { + _logger.LogWarning(exp, "Exception while deserializing workflow instance state: {InstanceId}. Reverting to default state", entity.Id); } - entity.WorkflowState = data; } diff --git a/src/modules/Elsa.Retention/Jobs/CleanupJob.cs b/src/modules/Elsa.Retention/Jobs/CleanupJob.cs index 764e21395..691963636 100644 --- a/src/modules/Elsa.Retention/Jobs/CleanupJob.cs +++ b/src/modules/Elsa.Retention/Jobs/CleanupJob.cs @@ -71,7 +71,14 @@ public class CleanupJob( await foreach (var entities in collector.GetRelatedEntitiesGeneric(page.Items).WithCancellation(cancellationToken)) { - await cleanupService.Cleanup(entities); + try + { + await cleanupService.Cleanup(entities); + } + catch (Exception exp) + { + _logger.LogError(exp, "Failed to clean up {Type} because exception thrown", collectorService.Key.Name); + } } }