From ed14a1e577411bb8b7bf7ba8b4e2424a18e6a5a2 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 18 Jul 2025 14:14:19 +0200 Subject: [PATCH] Refactor activity execution record serialization with snapshots (#6807) - Introduced `ActivityExecutionRecordSnapshot` for encapsulated serialized data. - Updated `DefaultActivityExecutionMapper` to build serialized snapshots. - Adjusted `ActivityExecutionLogStore` to persist pre-serialized snapshots. - Streamlined package version management with `MicrosoftVersion` property. --- Directory.Packages.props | 67 ++++++++++--------- src/modules/Elsa.Common/Entities/Entity.cs | 2 +- .../Runtime/ActivityExecutionLogStore.cs | 29 ++++---- .../Entities/ActivityExecutionRecord.cs | 9 +-- .../Models/ActivityExecutionRecordSnapshot.cs | 25 +++++++ .../DefaultActivityExecutionMapper.cs | 48 ++++++++----- 6 files changed, 108 insertions(+), 72 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs diff --git a/Directory.Packages.props b/Directory.Packages.props index 9a0667e87..8b65fbfaa 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -5,6 +5,7 @@ 3.5.0-preview.1092 + 9.0.7 @@ -34,7 +35,7 @@ - + @@ -118,8 +119,8 @@ - - + + @@ -127,36 +128,36 @@ - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/modules/Elsa.Common/Entities/Entity.cs b/src/modules/Elsa.Common/Entities/Entity.cs index 2a3d2cd6e..3da8462f6 100644 --- a/src/modules/Elsa.Common/Entities/Entity.cs +++ b/src/modules/Elsa.Common/Entities/Entity.cs @@ -8,7 +8,7 @@ public abstract class Entity /// /// Gets or sets the ID of this entity. /// - public string Id { get; set; } = default!; + public string Id { get; set; } = null!; /// /// Gets or sets the ID of the tenant that own this entity. diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs index a71d19e6a..92b684ac9 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs @@ -6,17 +6,13 @@ using Elsa.Common.Codecs; using Elsa.Common.Entities; using Elsa.Extensions; using Elsa.Workflows; -using Elsa.Workflows.Management.Options; using Elsa.Workflows.Runtime; using Elsa.Workflows.Runtime.Entities; -using Elsa.Workflows.Runtime.Extensions; using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.OrderDefinitions; using Elsa.Workflows.State; using JetBrains.Annotations; using Microsoft.EntityFrameworkCore; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; using Open.Linq.AsyncExtensions; namespace Elsa.EntityFrameworkCore.Modules.Runtime; @@ -29,8 +25,7 @@ public class EFCoreActivityExecutionStore( EntityStore store, ISafeSerializer safeSerializer, IPayloadSerializer payloadSerializer, - ICompressionCodecResolver compressionCodecResolver, - IOptions options) : IActivityExecutionStore + ICompressionCodecResolver compressionCodecResolver) : IActivityExecutionStore { /// public async Task SaveAsync(ActivityExecutionRecord record, CancellationToken cancellationToken = default) => await store.SaveAsync(record, OnSaveAsync, cancellationToken); @@ -87,15 +82,21 @@ public class EFCoreActivityExecutionStore( return await store.DeleteWhereAsync(queryable => Filter(queryable, filter), cancellationToken); } - private async ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken) + private ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken) { - dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = entity.SerializedActivityState; - dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = entity.SerializedActivityStateCompressionAlgorithm ?? nameof(None); - dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = entity.SerializedOutputs; - dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = entity.SerializedProperties; - dbContext.Entry(entity).Property("SerializedMetadata").CurrentValue = entity.SerializedMetadata; - dbContext.Entry(entity).Property("SerializedException").CurrentValue = entity.SerializedException; - dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = entity.SerializedPayload; + var snapshot = entity.SerializedSnapshot; + + if (snapshot is null) + return ValueTask.CompletedTask; + + dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = snapshot.SerializedActivityState; + dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = snapshot.SerializedActivityStateCompressionAlgorithm; + dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = snapshot.SerializedOutputs; + dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = snapshot.SerializedProperties; + dbContext.Entry(entity).Property("SerializedMetadata").CurrentValue = snapshot.SerializedMetadata; + dbContext.Entry(entity).Property("SerializedException").CurrentValue = snapshot.SerializedException; + dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = snapshot.SerializedPayload; + return ValueTask.CompletedTask; } [RequiresUnreferencedCode("Calls Elsa.EntityFrameworkCore.Modules.Runtime.EFCoreActivityExecutionStore.DeserializeActivityState(RuntimeElsaDbContext, ActivityExecutionRecord, CancellationToken)")] diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs index fac785904..2b6bd1e74 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs @@ -1,4 +1,5 @@ using System.ComponentModel.DataAnnotations.Schema; +using System.Text.Json.Serialization; using Elsa.Common; using Elsa.Common.Entities; using Elsa.Workflows.State; @@ -99,11 +100,5 @@ public partial class ActivityExecutionRecord : Entity, ILogRecord public partial class ActivityExecutionRecord { - [NotMapped] public string? SerializedActivityState { get; set; } - [NotMapped] public string? SerializedOutputs { get; set; } - [NotMapped] public string? SerializedProperties { get; set; } - [NotMapped] public string? SerializedPayload { get; set; } - [NotMapped] public string? SerializedMetadata { get; set; } - [NotMapped] public string? SerializedException { get; set; } - [NotMapped] public string? SerializedActivityStateCompressionAlgorithm { get; set; } + [NotMapped] [JsonIgnore] public ActivityExecutionRecordSnapshot? SerializedSnapshot { get; set; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs new file mode 100644 index 000000000..377850153 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs @@ -0,0 +1,25 @@ +namespace Elsa.Workflows.Runtime; + +public class ActivityExecutionRecordSnapshot +{ + public string Id { get; set; } = null!; + public string? TenantId { get; set; } + public string WorkflowInstanceId { get; set; } = null!; + public string ActivityId { get; set; } = null!; + public string ActivityNodeId { get; set; } = null!; + public string ActivityType { get; set; } = null!; + public int ActivityTypeVersion { get; set; } + public string? ActivityName { get; set; } + public DateTimeOffset StartedAt { get; set; } + public bool HasBookmarks { get; set; } + public ActivityStatus Status { get; set; } + public int AggregateFaultCount { get; set; } + public DateTimeOffset? CompletedAt { get; set; } + public string? SerializedActivityState { get; set; } + public string? SerializedOutputs { get; set; } + public string? SerializedProperties { get; set; } + public string? SerializedPayload { get; set; } + public string? SerializedMetadata { get; set; } + public string? SerializedException { get; set; } + public string? SerializedActivityStateCompressionAlgorithm { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs index ea3de58c2..223cd55f3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -11,12 +11,11 @@ namespace Elsa.Workflows.Runtime; /// public class DefaultActivityExecutionMapper( - ISafeSerializer safeSerializer, + ISafeSerializer safeSerializer, IPayloadSerializer payloadSerializer, ICompressionCodecResolver compressionCodecResolver, IOptions options) : IActivityExecutionMapper { - /// public async Task MapAsync(ActivityExecutionContext source) { @@ -50,22 +49,37 @@ public class DefaultActivityExecutionMapper( AggregateFaultCount = source.AggregateFaultCount, CompletedAt = source.CompletedAt }; - - record = record.SanitizeLogMessage(); - var compressionAlgorithm = options.Value.CompressionAlgorithm ?? nameof(None); - var serializedActivityState = record.ActivityState?.Count > 0 ? safeSerializer.Serialize(record.ActivityState) : null; - var compressedSerializedActivityState = serializedActivityState != null ? await compressionCodecResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : null; - var serializedProperties = record.Properties != null ? payloadSerializer.Serialize(record.Properties) : null; - var serializedMetadata = record.Metadata != null ? payloadSerializer.Serialize(record.Metadata) : null; - - record.SerializedActivityState = compressedSerializedActivityState; - record.SerializedActivityStateCompressionAlgorithm = compressionAlgorithm; - record.SerializedOutputs = record.Outputs?.Any() == true ? safeSerializer.Serialize(record.Outputs) : null; - record.SerializedProperties = serializedProperties; - record.SerializedMetadata = serializedMetadata; - record.SerializedException = record.Exception != null ? payloadSerializer.Serialize(record.Exception) : null; - record.SerializedPayload = record.Payload?.Any() == true ? payloadSerializer.Serialize(record.Payload) : null; + record = record.SanitizeLogMessage(); + var compressionAlgorithm = options.Value.CompressionAlgorithm ?? nameof(None); + var serializedActivityState = record.ActivityState?.Count > 0 ? safeSerializer.Serialize(record.ActivityState) : null; + var compressedSerializedActivityState = serializedActivityState != null ? await compressionCodecResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : null; + var serializedProperties = record.Properties != null ? payloadSerializer.Serialize(record.Properties) : null; + var serializedMetadata = record.Metadata != null ? payloadSerializer.Serialize(record.Metadata) : null; + record.SerializedSnapshot = new() + { + Id = record.Id, + TenantId = record.TenantId, + WorkflowInstanceId = record.WorkflowInstanceId, + ActivityId = record.ActivityId, + ActivityNodeId = record.ActivityNodeId, + ActivityType = record.ActivityType, + ActivityTypeVersion = record.ActivityTypeVersion, + ActivityName = record.ActivityName, + StartedAt = record.StartedAt, + HasBookmarks = record.HasBookmarks, + Status = record.Status, + AggregateFaultCount = record.AggregateFaultCount, + CompletedAt = record.CompletedAt, + SerializedActivityState = compressedSerializedActivityState, + SerializedActivityStateCompressionAlgorithm = compressionAlgorithm, + SerializedOutputs = record.Outputs?.Any() == true ? safeSerializer.Serialize(record.Outputs) : null, + SerializedProperties = serializedProperties, + SerializedMetadata = serializedMetadata, + SerializedException = record.Exception != null ? payloadSerializer.Serialize(record.Exception) : null, + SerializedPayload = record.Payload?.Any() == true ? payloadSerializer.Serialize(record.Payload) : null + }; + return record; }