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.
This commit is contained in:
parent
904284ce6b
commit
ed14a1e577
|
|
@ -5,6 +5,7 @@
|
|||
</PropertyGroup>
|
||||
<PropertyGroup>
|
||||
<ElsaStudioVersion>3.5.0-preview.1092</ElsaStudioVersion>
|
||||
<MicrosoftVersion>9.0.7</MicrosoftVersion>
|
||||
</PropertyGroup>
|
||||
<ItemGroup>
|
||||
<PackageVersion Include="Antlr4.Runtime.Standard" Version="4.13.1"/>
|
||||
|
|
@ -34,7 +35,7 @@
|
|||
<PackageVersion Include="DistributedLock.FileSystem" Version="1.0.3"/>
|
||||
<PackageVersion Include="DistributedLock.Postgres" Version="1.3.0"/>
|
||||
<PackageVersion Include="DistributedLock.Redis" Version="1.0.3"/>
|
||||
<PackageVersion Include="Elastic.Clients.Elasticsearch" Version="9.0.6"/>
|
||||
<PackageVersion Include="Elastic.Clients.Elasticsearch" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Elsa.Studio" Version="$(ElsaStudioVersion)"/>
|
||||
<PackageVersion Include="Elsa.Studio.Agents" Version="$(ElsaStudioVersion)"/>
|
||||
<PackageVersion Include="Elsa.Studio.Core.BlazorWasm" Version="$(ElsaStudioVersion)"/>
|
||||
|
|
@ -118,8 +119,8 @@
|
|||
<PackageVersion Include="System.Linq.Dynamic.Core" Version="1.6.5"/>
|
||||
<PackageVersion Include="System.Net.Http" Version="4.3.4"/>
|
||||
<PackageVersion Include="System.Text.RegularExpressions" Version="4.3.1"/>
|
||||
<PackageVersion Include="System.Formats.Asn1" Version="9.0.6"/>
|
||||
<PackageVersion Include="System.Text.Json" Version="9.0.6"/>
|
||||
<PackageVersion Include="System.Formats.Asn1" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="System.Text.Json" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Testcontainers" Version="4.5.0"/>
|
||||
<PackageVersion Include="Testcontainers.PostgreSql" Version="4.5.0"/>
|
||||
<PackageVersion Include="Testcontainers.RabbitMq" Version="4.5.0"/>
|
||||
|
|
@ -127,36 +128,36 @@
|
|||
<PackageVersion Include="ThrottleDebounce" Version="2.0.1"/>
|
||||
<PackageVersion Include="WebhooksCore" Version="0.0.1"/>
|
||||
<PackageVersion Include="Yarp.ReverseProxy" Version="2.3.0"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Authorization" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.DevServer" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.Server" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.DataProtection.Abstractions" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Mvc.Testing" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Data.Sqlite" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Data.Sqlite.Core" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.Design" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.Relational" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.Sqlite" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.SqlServer" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Caching.Abstractions" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Caching.Memory" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Configuration" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Configuration.Abstractions" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Configuration.Json" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.DependencyModel" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Hosting.Abstractions" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Http" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Http.Polly" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Logging" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Logging.Abstractions" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Logging.Console" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Options" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Options.ConfigurationExtensions" Version="9.0.6"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Authorization" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.DevServer" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Components.WebAssembly.Server" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.DataProtection.Abstractions" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.AspNetCore.Mvc.Testing" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Data.Sqlite" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Data.Sqlite.Core" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.Design" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.Relational" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.Sqlite" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.EntityFrameworkCore.SqlServer" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Caching.Abstractions" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Caching.Memory" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Configuration" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Configuration.Abstractions" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Configuration.Json" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.DependencyModel" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Hosting.Abstractions" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Http" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Http.Polly" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Logging" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Logging.Abstractions" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Logging.Console" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Options" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Options.ConfigurationExtensions" Version="$(MicrosoftVersion)"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Http.Resilience" Version="9.6.0"/>
|
||||
<PackageVersion Include="Microsoft.Extensions.Resilience" Version="9.6.0"/>
|
||||
<PackageVersion Include="MySql.Data" Version="9.3.0"/>
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ public abstract class Entity
|
|||
/// <summary>
|
||||
/// Gets or sets the ID of this entity.
|
||||
/// </summary>
|
||||
public string Id { get; set; } = default!;
|
||||
public string Id { get; set; } = null!;
|
||||
|
||||
/// <summary>
|
||||
/// Gets or sets the ID of the tenant that own this entity.
|
||||
|
|
|
|||
|
|
@ -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<RuntimeElsaDbContext, ActivityExecutionRecord> store,
|
||||
ISafeSerializer safeSerializer,
|
||||
IPayloadSerializer payloadSerializer,
|
||||
ICompressionCodecResolver compressionCodecResolver,
|
||||
IOptions<ManagementOptions> options) : IActivityExecutionStore
|
||||
ICompressionCodecResolver compressionCodecResolver) : IActivityExecutionStore
|
||||
{
|
||||
/// <inheritdoc />
|
||||
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)")]
|
||||
|
|
|
|||
|
|
@ -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; }
|
||||
}
|
||||
|
|
@ -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; }
|
||||
}
|
||||
|
|
@ -11,12 +11,11 @@ namespace Elsa.Workflows.Runtime;
|
|||
|
||||
/// <inheritdoc />
|
||||
public class DefaultActivityExecutionMapper(
|
||||
ISafeSerializer safeSerializer,
|
||||
ISafeSerializer safeSerializer,
|
||||
IPayloadSerializer payloadSerializer,
|
||||
ICompressionCodecResolver compressionCodecResolver,
|
||||
IOptions<ManagementOptions> options) : IActivityExecutionMapper
|
||||
{
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<ActivityExecutionRecord> 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;
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue