Introduce activity execution record capturing and serialization improvements (#6800)

* Introduce activity execution record capturing and serialization improvements

- Added middleware for capturing activity execution records during workflow execution.
- Introduced async mapping in `DefaultActivityExecutionMapper` with additional serialization support.
- Enhanced `ActivityExecutionRecord` with new serialized properties for efficient storage.
- Updated extensions to include `UseActivityExecutionLogCapturing`.
- Simplified logging persistence by leveraging pre-serialized values in `ActivityExecutionLogStore`.

* Update src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

* Update src/modules/Elsa.Workflows.Runtime/Middleware/Activities/CaptureActivityExecutionRecordMiddleware.cs

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

* Apply suggestion from @Copilot

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

---------

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
Sipke Schoorstra 2025-07-16 18:05:32 +02:00 committed by GitHub
parent 51fb7bfaef
commit adbea90ccd
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
11 changed files with 110 additions and 50 deletions

View file

@ -89,20 +89,13 @@ public class EFCoreActivityExecutionStore(
private async ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken)
{
entity = entity.SanitizeLogMessage();
var compressionAlgorithm = options.Value.CompressionAlgorithm ?? nameof(None);
var serializedActivityState = entity.ActivityState?.Count > 0 ? safeSerializer.Serialize(entity.ActivityState) : null;
var compressedSerializedActivityState = serializedActivityState != null ? await compressionCodecResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : null;
var serializedProperties = entity.Properties != null ? payloadSerializer.Serialize(entity.Properties) : null;
var serializedMetadata = entity.Metadata != null ? payloadSerializer.Serialize(entity.Metadata) : null;
dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = compressedSerializedActivityState;
dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = compressionAlgorithm;
dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = entity.Outputs?.Any() == true ? safeSerializer.Serialize(entity.Outputs) : null;
dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = serializedProperties;
dbContext.Entry(entity).Property("SerializedMetadata").CurrentValue = serializedMetadata;
dbContext.Entry(entity).Property("SerializedException").CurrentValue = entity.Exception != null ? payloadSerializer.Serialize(entity.Exception) : null;
dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = entity.Payload?.Any() == true ? payloadSerializer.Serialize(entity.Payload) : null;
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;
}
[RequiresUnreferencedCode("Calls Elsa.EntityFrameworkCore.Modules.Runtime.EFCoreActivityExecutionStore.DeserializeActivityState(RuntimeElsaDbContext, ActivityExecutionRecord, CancellationToken)")]
@ -129,7 +122,7 @@ public class EFCoreActivityExecutionStore(
var compressionAlgorithm = (string?)dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue ?? nameof(None);
var compressionStrategy = compressionCodecResolver.Resolve(compressionAlgorithm);
json = await compressionStrategy.DecompressAsync(json, cancellationToken);
var dictionary = JsonSerializer.Deserialize<IDictionary<string, object?>>(json);
var dictionary = safeSerializer.Deserialize<IDictionary<string, object?>?>(json);
return dictionary?.ToDictionary(x => x.Key, x => x.Value);
}

View file

@ -82,14 +82,6 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I
// Invoke next middleware.
await next(context);
// // If the activity created any bookmarks, copy them into the workflow execution context.
// if (context.Bookmarks.Any())
// {
// // Store bookmarks.
// workflowExecutionContext.Bookmarks.AddRange(context.Bookmarks);
// logger.LogDebug("Added {BookmarkCount} bookmarks to the workflow execution context", context.Bookmarks.Count);
// }
// Conditionally commit the workflow state.
if (ShouldCommit(context, ActivityLifetimeEvent.ActivityExecuted))
await context.WorkflowExecutionContext.CommitAsync();

View file

@ -10,11 +10,5 @@ public interface IActivityExecutionMapper
/// <summary>
/// Maps an activity execution context to an activity execution record.
/// </summary>
ActivityExecutionRecord Map(ActivityExecutionContext source);
/// <summary>
/// Maps an activity execution context to an activity execution record.
/// </summary>
[Obsolete( "Use Map instead.", error: false)]
Task<ActivityExecutionRecord> MapAsync(ActivityExecutionContext source);
}

View file

@ -1,3 +1,4 @@
using System.ComponentModel.DataAnnotations.Schema;
using Elsa.Common;
using Elsa.Common.Entities;
using Elsa.Workflows.State;
@ -7,18 +8,18 @@ namespace Elsa.Workflows.Runtime.Entities;
/// <summary>
/// Represents a single activity execution of an activity instance.
/// </summary>
public class ActivityExecutionRecord : Entity, ILogRecord
public partial class ActivityExecutionRecord : Entity, ILogRecord
{
/// <summary>
/// Gets or sets the workflow instance ID.
/// </summary>
public string WorkflowInstanceId { get; set; } = null!;
/// <summary>
/// Gets or sets the activity ID.
/// </summary>
public string ActivityId { get; set; } = null!;
/// <summary>
/// Gets or sets the activity node ID.
/// </summary>
@ -38,17 +39,17 @@ public class ActivityExecutionRecord : Entity, ILogRecord
/// The name of the activity.
/// </summary>
public string? ActivityName { get; set; }
/// <summary>
/// The state of the activity at the time this record is created or last updated.
/// </summary>
public IDictionary<string, object?>? ActivityState { get; set; }
/// <summary>
/// Any additional payload associated with the log record.
/// </summary>
public IDictionary<string, object>? Payload { get; set; }
/// <summary>
/// Any outputs provided by the activity.
/// </summary>
@ -58,7 +59,7 @@ public class ActivityExecutionRecord : Entity, ILogRecord
/// Any properties provided by the activity.
/// </summary>
public IDictionary<string, object>? Properties { get; set; }
/// <summary>
/// Lightweight metadata associated with the activity execution.
/// This information will be retained as part of the activity execution summary record.
@ -79,7 +80,7 @@ public class ActivityExecutionRecord : Entity, ILogRecord
/// Gets or sets whether the activity has any bookmarks.
/// </summary>
public bool HasBookmarks { get; set; }
/// <summary>
/// Gets or sets the status of the activity.
/// </summary>
@ -89,9 +90,20 @@ public class ActivityExecutionRecord : Entity, ILogRecord
/// Gets or sets the aggregated count of faults encountered during the execution of the activity instance and its descendants.
/// </summary>
public int AggregateFaultCount { get; set; }
/// <summary>
/// Gets or sets the time at which the activity execution completed.
/// </summary>
public DateTimeOffset? CompletedAt { get; set; }
}
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; }
}

View file

@ -7,7 +7,7 @@ using Elsa.Workflows.Runtime.Stimuli;
// ReSharper disable once CheckNamespace
namespace Elsa.Extensions;
public static class ActivityExecutionContextExtensions
public static class ActivityExecutionContextEventExtensions
{
/// <summary>
/// Suspends the current activity's execution and waits for a specified event to occur before continuing.

View file

@ -0,0 +1,23 @@
using Elsa.Workflows;
using Elsa.Workflows.Runtime;
using Elsa.Workflows.Runtime.Entities;
// ReSharper disable once CheckNamespace
namespace Elsa.Extensions;
public static class ActivityExecutionContextRecordExtensions
{
private const string ActivityExecutionRecordKey = "CapturedActivityExecutionRecord";
public static async Task CaptureActivityExecutionRecordAsync(this ActivityExecutionContext context)
{
var mapper = context.GetRequiredService<IActivityExecutionMapper>();
var record = await mapper.MapAsync(context);
context.TransientProperties[ActivityExecutionRecordKey] = record;
}
public static ActivityExecutionRecord? GetCapturedActivityExecutionRecord(this ActivityExecutionContext context)
{
return context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var record) ? (ActivityExecutionRecord?)record : null;
}
}

View file

@ -34,6 +34,7 @@ public static class PipelineWorkflowsFeatureExtensions
.UseExecutionLogging()
.UseNotifications()
.UseLogPersistenceModeEvaluation()
.UseActivityExecutionLogCapturing()
.UseBackgroundActivityInvoker();
configurePipeline?.Invoke(pipeline);

View file

@ -20,4 +20,12 @@ public static class ActivityExecutionPipelineBuilderExtensions
/// Installs the <see cref="EvaluateLogPersistenceModesMiddleware"/> which evaluates log persistence modes during activity execution.
/// </summary>
public static IActivityExecutionPipelineBuilder UseLogPersistenceModeEvaluation(this IActivityExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<EvaluateLogPersistenceModesMiddleware>();
/// <summary>
/// Installs the <see cref="CaptureActivityExecutionRecordMiddleware"/> into the activity execution pipeline to capture and map activity execution details.
/// </summary>
public static IActivityExecutionPipelineBuilder UseActivityExecutionLogCapturing(this IActivityExecutionPipelineBuilder pipelineBuilder)
{
return pipelineBuilder.UseMiddleware<CaptureActivityExecutionRecordMiddleware>();
}
}

View file

@ -0,0 +1,13 @@
using Elsa.Extensions;
using Elsa.Workflows.Pipelines.ActivityExecution;
namespace Elsa.Workflows.Runtime.Middleware.Activities;
public class CaptureActivityExecutionRecordMiddleware(ActivityMiddlewareDelegate next) : IActivityExecutionMiddleware
{
public async ValueTask InvokeAsync(ActivityExecutionContext context)
{
await next(context);
await context.CaptureActivityExecutionRecordAsync();
}
}

View file

@ -1,13 +1,24 @@
using Elsa.Common;
using Elsa.Common.Codecs;
using Elsa.Workflows.LogPersistence;
using Elsa.Workflows.Management.Options;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Extensions;
using Elsa.Workflows.State;
using Microsoft.Extensions.Options;
namespace Elsa.Workflows.Runtime;
/// <inheritdoc />
public class DefaultActivityExecutionMapper() : IActivityExecutionMapper
public class DefaultActivityExecutionMapper(
ISafeSerializer safeSerializer,
IPayloadSerializer payloadSerializer,
ICompressionCodecResolver compressionCodecResolver,
IOptions<ManagementOptions> options) : IActivityExecutionMapper
{
public ActivityExecutionRecord Map(ActivityExecutionContext source)
/// <inheritdoc />
public async Task<ActivityExecutionRecord> MapAsync(ActivityExecutionContext source)
{
var outputs = source.GetOutputs();
var inputs = source.GetInputs();
@ -16,8 +27,9 @@ public class DefaultActivityExecutionMapper() : IActivityExecutionMapper
var persistableOutputs = GetPersistableInputOutput(outputs, persistenceMap.Outputs);
var persistableProperties = GetPersistableDictionary(source.Properties!, persistenceMap.InternalState);
var persistableJournalData = GetPersistableDictionary(source.JournalData!, persistenceMap.InternalState);
var cancellationToken = source.CancellationToken;
return new()
var record = new ActivityExecutionRecord
{
Id = source.Id,
ActivityId = source.Activity.Id,
@ -38,15 +50,26 @@ public class DefaultActivityExecutionMapper() : IActivityExecutionMapper
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;
return record;
}
/// <inheritdoc />
public Task<ActivityExecutionRecord> MapAsync(ActivityExecutionContext source)
{
return Task.FromResult(Map(source));
}
private IDictionary<string, object?> GetPersistableInputOutput(IDictionary<string, object> state, IDictionary<string, LogPersistenceMode> map)
private IDictionary<string, object?> GetPersistableInputOutput(IDictionary<string, object> state, IDictionary<string, LogPersistenceMode> map, bool deepCopy = false)
{
var result = new Dictionary<string, object?>();
foreach (var stateEntry in state)

View file

@ -1,3 +1,4 @@
using Elsa.Extensions;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Notifications;
@ -22,7 +23,7 @@ public class StoreActivityExecutionLogSink(
if (activityExecutionContexts.Count == 0)
return;
var records = activityExecutionContexts.Select(mapper.Map).ToList();
var records = await Task.WhenAll(activityExecutionContexts.Select(async x => x.GetCapturedActivityExecutionRecord() ?? await mapper.MapAsync(x)));
await activityExecutionStore.SaveManyAsync(records, cancellationToken);
// Untaint activity execution contexts.