From 73482a8e4c3d3590d152623ec8cda6ed98c84492 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 23 Aug 2023 13:34:30 +0200 Subject: [PATCH] Enhanced Logging for Activity Outcomes and Outputs (#4352) * Add result to logs * Include outcomes and outputs with activity execution log * Include outputs with completed activity events --- Elsa.sln | 2 +- ....cs => 20230823105545_Initial.Designer.cs} | 8 ++- ...0_Initial.cs => 20230823105545_Initial.cs} | 4 ++ .../RuntimeElsaDbContextModelSnapshot.cs | 6 +++ ....cs => 20230823105601_Initial.Designer.cs} | 8 ++- ...9_Initial.cs => 20230823105601_Initial.cs} | 4 +- .../RuntimeElsaDbContextModelSnapshot.cs | 6 +++ ....cs => 20230823105550_Initial.Designer.cs} | 8 ++- ...5_Initial.cs => 20230823105550_Initial.cs} | 4 +- .../RuntimeElsaDbContextModelSnapshot.cs | 6 +++ ....cs => 20230823105555_Initial.Designer.cs} | 8 ++- ...7_Initial.cs => 20230823105555_Initial.cs} | 4 +- .../RuntimeElsaDbContextModelSnapshot.cs | 6 +++ .../Runtime/ActivityExecutionLogStore.cs | 47 +++++++++++------- .../Modules/Runtime/Configurations.cs | 4 ++ .../ActivityExecutionContextExtensions.cs | 28 ++++++++++- .../Entities/ActivityExecutionRecord.cs | 10 ++++ .../PersistActivityExecutionLogMiddleware.cs | 49 +++++++++++++------ ...rations-runtime.sh => update-migrations.sh | 0 19 files changed, 170 insertions(+), 42 deletions(-) rename src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/{20230821204450_Initial.Designer.cs => 20230823105545_Initial.Designer.cs} (98%) rename src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/{20230821204450_Initial.cs => 20230823105545_Initial.cs} (98%) rename src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/{20230821204529_Initial.Designer.cs => 20230823105601_Initial.Designer.cs} (98%) rename src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/{20230821204529_Initial.cs => 20230823105601_Initial.cs} (99%) rename src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/{20230821204455_Initial.Designer.cs => 20230823105550_Initial.Designer.cs} (98%) rename src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/{20230821204455_Initial.cs => 20230823105550_Initial.cs} (98%) rename src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/{20230821204507_Initial.Designer.cs => 20230823105555_Initial.Designer.cs} (98%) rename src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/{20230821204507_Initial.cs => 20230823105555_Initial.cs} (98%) rename update-migrations-runtime.sh => update-migrations.sh (100%) diff --git a/Elsa.sln b/Elsa.sln index 0b2b6cb9b..baf55f4ac 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -19,7 +19,7 @@ ProjectSection(SolutionItems) = preProject build-and-run-all-in-one-web-docker.sh = build-and-run-all-in-one-web-docker.sh NuGet.Config = NuGet.Config .github\workflows\npm-packages.yml = .github\workflows\npm-packages.yml - update-migrations-runtime.sh = update-migrations-runtime.sh + update-migrations.sh = update-migrations.sh .github\workflows\pr-body-generator.yml = .github\workflows\pr-body-generator.yml EndProjectSection EndProject diff --git a/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230821204450_Initial.Designer.cs b/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230823105545_Initial.Designer.cs similarity index 98% rename from src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230821204450_Initial.Designer.cs rename to src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230823105545_Initial.Designer.cs index 02a6d87e2..618a08a1a 100644 --- a/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230821204450_Initial.Designer.cs +++ b/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230823105545_Initial.Designer.cs @@ -11,7 +11,7 @@ using Microsoft.EntityFrameworkCore.Storage.ValueConversion; namespace Elsa.EntityFrameworkCore.MySql.Migrations.Runtime { [DbContext(typeof(RuntimeElsaDbContext))] - [Migration("20230821204450_Initial")] + [Migration("20230823105545_Initial")] partial class Initial { /// @@ -125,6 +125,12 @@ namespace Elsa.EntityFrameworkCore.MySql.Migrations.Runtime b.Property("SerializedException") .HasColumnType("longtext"); + b.Property("SerializedOutputs") + .HasColumnType("longtext"); + + b.Property("SerializedPayload") + .HasColumnType("longtext"); + b.Property("StartedAt") .HasColumnType("datetime(6)"); diff --git a/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230821204450_Initial.cs b/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230823105545_Initial.cs similarity index 98% rename from src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230821204450_Initial.cs rename to src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230823105545_Initial.cs index f024f9926..57853edac 100644 --- a/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230821204450_Initial.cs +++ b/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/20230823105545_Initial.cs @@ -40,6 +40,10 @@ namespace Elsa.EntityFrameworkCore.MySql.Migrations.Runtime SerializedActivityState = table.Column(type: "longtext", nullable: true) .Annotation("MySql:CharSet", "utf8mb4"), SerializedException = table.Column(type: "longtext", nullable: true) + .Annotation("MySql:CharSet", "utf8mb4"), + SerializedOutputs = table.Column(type: "longtext", nullable: true) + .Annotation("MySql:CharSet", "utf8mb4"), + SerializedPayload = table.Column(type: "longtext", nullable: true) .Annotation("MySql:CharSet", "utf8mb4") }, constraints: table => diff --git a/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs b/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs index 9d8ff4eaf..7cbff57a6 100644 --- a/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs +++ b/src/modules/Elsa.EntityFrameworkCore.MySql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs @@ -122,6 +122,12 @@ namespace Elsa.EntityFrameworkCore.MySql.Migrations.Runtime b.Property("SerializedException") .HasColumnType("longtext"); + b.Property("SerializedOutputs") + .HasColumnType("longtext"); + + b.Property("SerializedPayload") + .HasColumnType("longtext"); + b.Property("StartedAt") .HasColumnType("datetime(6)"); diff --git a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230821204529_Initial.Designer.cs b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230823105601_Initial.Designer.cs similarity index 98% rename from src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230821204529_Initial.Designer.cs rename to src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230823105601_Initial.Designer.cs index 1a62de5fb..ae1c8a412 100644 --- a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230821204529_Initial.Designer.cs +++ b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230823105601_Initial.Designer.cs @@ -12,7 +12,7 @@ using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; namespace Elsa.EntityFrameworkCore.PostgreSql.Migrations.Runtime { [DbContext(typeof(RuntimeElsaDbContext))] - [Migration("20230821204529_Initial")] + [Migration("20230823105601_Initial")] partial class Initial { /// @@ -128,6 +128,12 @@ namespace Elsa.EntityFrameworkCore.PostgreSql.Migrations.Runtime b.Property("SerializedException") .HasColumnType("text"); + b.Property("SerializedOutputs") + .HasColumnType("text"); + + b.Property("SerializedPayload") + .HasColumnType("text"); + b.Property("StartedAt") .HasColumnType("timestamp with time zone"); diff --git a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230821204529_Initial.cs b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230823105601_Initial.cs similarity index 99% rename from src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230821204529_Initial.cs rename to src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230823105601_Initial.cs index b70ab87b7..58e185b41 100644 --- a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230821204529_Initial.cs +++ b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/20230823105601_Initial.cs @@ -30,7 +30,9 @@ namespace Elsa.EntityFrameworkCore.PostgreSql.Migrations.Runtime Status = table.Column(type: "integer", nullable: false), CompletedAt = table.Column(type: "timestamp with time zone", nullable: true), SerializedActivityState = table.Column(type: "text", nullable: true), - SerializedException = table.Column(type: "text", nullable: true) + SerializedException = table.Column(type: "text", nullable: true), + SerializedOutputs = table.Column(type: "text", nullable: true), + SerializedPayload = table.Column(type: "text", nullable: true) }, constraints: table => { diff --git a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs index e6f84031c..cfec1da8d 100644 --- a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs +++ b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs @@ -125,6 +125,12 @@ namespace Elsa.EntityFrameworkCore.PostgreSql.Migrations.Runtime b.Property("SerializedException") .HasColumnType("text"); + b.Property("SerializedOutputs") + .HasColumnType("text"); + + b.Property("SerializedPayload") + .HasColumnType("text"); + b.Property("StartedAt") .HasColumnType("timestamp with time zone"); diff --git a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230821204455_Initial.Designer.cs b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230823105550_Initial.Designer.cs similarity index 98% rename from src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230821204455_Initial.Designer.cs rename to src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230823105550_Initial.Designer.cs index cb89f56c0..db17949ba 100644 --- a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230821204455_Initial.Designer.cs +++ b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230823105550_Initial.Designer.cs @@ -12,7 +12,7 @@ using Microsoft.EntityFrameworkCore.Storage.ValueConversion; namespace Elsa.EntityFrameworkCore.SqlServer.Migrations.Runtime { [DbContext(typeof(RuntimeElsaDbContext))] - [Migration("20230821204455_Initial")] + [Migration("20230823105550_Initial")] partial class Initial { /// @@ -128,6 +128,12 @@ namespace Elsa.EntityFrameworkCore.SqlServer.Migrations.Runtime b.Property("SerializedException") .HasColumnType("nvarchar(max)"); + b.Property("SerializedOutputs") + .HasColumnType("nvarchar(max)"); + + b.Property("SerializedPayload") + .HasColumnType("nvarchar(max)"); + b.Property("StartedAt") .HasColumnType("datetimeoffset"); diff --git a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230821204455_Initial.cs b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230823105550_Initial.cs similarity index 98% rename from src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230821204455_Initial.cs rename to src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230823105550_Initial.cs index 1a3181d31..1d717d45a 100644 --- a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230821204455_Initial.cs +++ b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/20230823105550_Initial.cs @@ -30,7 +30,9 @@ namespace Elsa.EntityFrameworkCore.SqlServer.Migrations.Runtime Status = table.Column(type: "int", nullable: false), CompletedAt = table.Column(type: "datetimeoffset", nullable: true), SerializedActivityState = table.Column(type: "nvarchar(max)", nullable: true), - SerializedException = table.Column(type: "nvarchar(max)", nullable: true) + SerializedException = table.Column(type: "nvarchar(max)", nullable: true), + SerializedOutputs = table.Column(type: "nvarchar(max)", nullable: true), + SerializedPayload = table.Column(type: "nvarchar(max)", nullable: true) }, constraints: table => { diff --git a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs index 6d2382a23..d94741068 100644 --- a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs +++ b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs @@ -125,6 +125,12 @@ namespace Elsa.EntityFrameworkCore.SqlServer.Migrations.Runtime b.Property("SerializedException") .HasColumnType("nvarchar(max)"); + b.Property("SerializedOutputs") + .HasColumnType("nvarchar(max)"); + + b.Property("SerializedPayload") + .HasColumnType("nvarchar(max)"); + b.Property("StartedAt") .HasColumnType("datetimeoffset"); diff --git a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230821204507_Initial.Designer.cs b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230823105555_Initial.Designer.cs similarity index 98% rename from src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230821204507_Initial.Designer.cs rename to src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230823105555_Initial.Designer.cs index 5f5f75265..d9c5e33aa 100644 --- a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230821204507_Initial.Designer.cs +++ b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230823105555_Initial.Designer.cs @@ -10,7 +10,7 @@ using Microsoft.EntityFrameworkCore.Storage.ValueConversion; namespace Elsa.EntityFrameworkCore.Sqlite.Migrations.Runtime { [DbContext(typeof(RuntimeElsaDbContext))] - [Migration("20230821204507_Initial")] + [Migration("20230823105555_Initial")] partial class Initial { /// @@ -123,6 +123,12 @@ namespace Elsa.EntityFrameworkCore.Sqlite.Migrations.Runtime b.Property("SerializedException") .HasColumnType("TEXT"); + b.Property("SerializedOutputs") + .HasColumnType("TEXT"); + + b.Property("SerializedPayload") + .HasColumnType("TEXT"); + b.Property("StartedAt") .IsRequired() .HasColumnType("TEXT"); diff --git a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230821204507_Initial.cs b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230823105555_Initial.cs similarity index 98% rename from src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230821204507_Initial.cs rename to src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230823105555_Initial.cs index cbfb07c43..a8c1bd16a 100644 --- a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230821204507_Initial.cs +++ b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/20230823105555_Initial.cs @@ -25,7 +25,9 @@ namespace Elsa.EntityFrameworkCore.Sqlite.Migrations.Runtime Status = table.Column(type: "INTEGER", nullable: false), CompletedAt = table.Column(type: "TEXT", nullable: true), SerializedActivityState = table.Column(type: "TEXT", nullable: true), - SerializedException = table.Column(type: "TEXT", nullable: true) + SerializedException = table.Column(type: "TEXT", nullable: true), + SerializedOutputs = table.Column(type: "TEXT", nullable: true), + SerializedPayload = table.Column(type: "TEXT", nullable: true) }, constraints: table => { diff --git a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs index 7997d03aa..3068a499d 100644 --- a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs +++ b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Migrations/Runtime/RuntimeElsaDbContextModelSnapshot.cs @@ -120,6 +120,12 @@ namespace Elsa.EntityFrameworkCore.Sqlite.Migrations.Runtime b.Property("SerializedException") .HasColumnType("TEXT"); + b.Property("SerializedOutputs") + .HasColumnType("TEXT"); + + b.Property("SerializedPayload") + .HasColumnType("TEXT"); + b.Property("StartedAt") .IsRequired() .HasColumnType("TEXT"); diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs index 96e98e4ed..4e596f3e7 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs @@ -38,46 +38,57 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore /// public async Task> FindManyAsync(ActivityExecutionRecordFilter filter, ActivityExecutionRecordOrder order, CancellationToken cancellationToken = default) => - await _store.QueryAsync(queryable => Filter(queryable, filter).OrderBy(order), OnLoadAsync, cancellationToken).ToList(); + await _store.QueryAsync(queryable => EFCoreActivityExecutionStore.Filter(queryable, filter).OrderBy(order), OnLoadAsync, cancellationToken).ToList(); /// public async Task> FindManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => - await _store.QueryAsync(queryable => Filter(queryable, filter), OnLoadAsync, cancellationToken).ToList(); + await _store.QueryAsync(queryable => EFCoreActivityExecutionStore.Filter(queryable, filter), OnLoadAsync, cancellationToken).ToList(); /// - public async Task CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.CountAsync(queryable => Filter(queryable, filter), cancellationToken); + public async Task CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.CountAsync(queryable => EFCoreActivityExecutionStore.Filter(queryable, filter), cancellationToken); /// - public async Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.DeleteWhereAsync(queryable => Filter(queryable, filter), cancellationToken); + public async Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.DeleteWhereAsync(queryable => EFCoreActivityExecutionStore.Filter(queryable, filter), cancellationToken); private async ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken) { dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = entity.ActivityState != null ? (await _activityStateSerializer.SerializeAsync(entity.ActivityState, cancellationToken)).ToString() : default; + dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = entity.ActivityState != null ? (await _activityStateSerializer.SerializeAsync(entity.Outputs, cancellationToken)).ToString() : default; dbContext.Entry(entity).Property("SerializedException").CurrentValue = entity.Exception != null ? _payloadSerializer.Serialize(entity.Exception) : default; + dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = entity.Payload != null ? _payloadSerializer.Serialize(entity.Payload) : default; } - private async ValueTask OnLoadAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord? entity, CancellationToken cancellationToken) + private ValueTask OnLoadAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord? entity, CancellationToken cancellationToken) { if (entity is null) - return; + return ValueTask.CompletedTask; - entity.ActivityState = await LoadActivityState(dbContext, entity); - entity.Exception = await LoadException(dbContext, entity); + entity.ActivityState = DeserializeActivityState(dbContext, entity); + entity.Outputs = Deserialize>(dbContext, entity, "SerializedOutputs"); + entity.Exception = DeserializePayload(dbContext, entity, "SerializedException"); + entity.Payload = DeserializePayload>(dbContext, entity, "SerializedPayload"); + return ValueTask.CompletedTask; } - private ValueTask?> LoadActivityState(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity) + private IDictionary? DeserializeActivityState(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity) { - var json = dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue; - var dictionary = !string.IsNullOrEmpty(json) ? JsonSerializer.Deserialize>(json) : default; - return ValueTask.FromResult?>(dictionary?.ToDictionary(x => x.Key, x => (object)x.Value)); + var dictionary = Deserialize>(dbContext, entity, "SerializedActivityState"); + return dictionary?.ToDictionary(x => x.Key, x => (object)x.Value); } - - private ValueTask LoadException(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity) + + private T? Deserialize(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, string propertyName) { - var json = dbContext.Entry(entity).Property("SerializedException").CurrentValue; - var exceptionState = !string.IsNullOrEmpty(json) ? _payloadSerializer.Deserialize(json) : default; - return ValueTask.FromResult(exceptionState); + var json = dbContext.Entry(entity).Property(propertyName).CurrentValue; + var value = !string.IsNullOrEmpty(json) ? JsonSerializer.Deserialize(json) : default; + return value; } - private IQueryable Filter(IQueryable queryable, ActivityExecutionRecordFilter filter) => filter.Apply(queryable); + private T? DeserializePayload(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, string propertyName) + { + var json = dbContext.Entry(entity).Property(propertyName).CurrentValue; + var payload = !string.IsNullOrEmpty(json) ? _payloadSerializer.Deserialize(json) : default; + return payload; + } + + private static IQueryable Filter(IQueryable queryable, ActivityExecutionRecordFilter filter) => filter.Apply(queryable); } \ No newline at end of file diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/Configurations.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/Configurations.cs index eab3cf262..0e3c2d3ce 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/Configurations.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/Configurations.cs @@ -82,8 +82,12 @@ public class Configurations : { builder.Ignore(x => x.ActivityState); builder.Ignore(x => x.Exception); + builder.Ignore(x => x.Payload); + builder.Ignore(x => x.Outputs); builder.Property("SerializedActivityState"); builder.Property("SerializedException"); + builder.Property("SerializedPayload"); + builder.Property("SerializedOutputs"); builder.HasIndex(x => x.WorkflowInstanceId).HasDatabaseName($"IX_{nameof(ActivityExecutionRecord)}_{nameof(ActivityExecutionRecord.WorkflowInstanceId)}"); builder.HasIndex(x => x.ActivityId).HasDatabaseName($"IX_{nameof(ActivityExecutionRecord)}_{nameof(ActivityExecutionRecord.ActivityId)}"); diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs index c44703071..3b3b117ed 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs @@ -338,12 +338,36 @@ public static class ActivityExecutionContextExtensions public static async ValueTask CompleteActivityAsync(this ActivityExecutionContext context, object? result = default) { // If the activity is already completed, do nothing. - if(context.Status == ActivityStatus.Completed) + if (context.Status == ActivityStatus.Completed) return; - + // Mark the activity as complete. context.Status = ActivityStatus.Completed; + // Record the outcomes, if any. + if (result is Outcomes outcomes) + context.JournalData["Outcomes"] = outcomes.Names; + + // Record the output, if any. + var activity = context.Activity; + var expressionExecutionContext = context.ExpressionExecutionContext; + var activityDescriptor = context.ActivityDescriptor; + var outputDescriptors = activityDescriptor.Outputs; + var outputs = outputDescriptors.ToDictionary(x => x.Name, x => activity.GetOutput(expressionExecutionContext, x.Name)!); + var serializer = context.GetRequiredService(); + + foreach (var output in outputs) + { + var outputName = output.Key; + var outputValue = output.Value; + + if (outputValue == null!) + continue; + + var serializedOutputValue = await serializer.SerializeAsync(outputValue, context.CancellationToken); + context.JournalData[outputName] = serializedOutputValue; + } + // Add an execution log entry. context.AddExecutionLogEntry("Completed", payload: context.JournalData, includeActivityState: true); diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs index 9e167843f..cb5e330c6 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs @@ -38,6 +38,16 @@ public class ActivityExecutionRecord : Entity /// The state of the activity at the time this record is created or last updated. /// public IDictionary? ActivityState { get; set; } + + /// + /// Any additional payload associated with the log record. + /// + public IDictionary? Payload { get; set; } + + /// + /// Any outputs provided by the activity. + /// + public IDictionary? Outputs { get; set; } /// /// Gets or sets the exception that occurred during the activity execution. diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs index ee7e65dc7..ac4f2ee52 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs @@ -1,5 +1,7 @@ +using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Core; +using Elsa.Workflows.Core.Activities.Flowchart.Models; using Elsa.Workflows.Core.Pipelines.WorkflowExecution; using Elsa.Workflows.Core.State; using Elsa.Workflows.Runtime.Contracts; @@ -31,22 +33,41 @@ public class PersistActivityExecutionLogMiddleware : WorkflowExecutionMiddleware // Get all activity execution contexts. var activityExecutionContexts = context.ActivityExecutionContexts; - + // Persist activity execution entries. - var entries = activityExecutionContexts.Select(x => new ActivityExecutionRecord + var entries = activityExecutionContexts.Select(activityExecutionContext => { - Id = x.Id, - ActivityId = x.Activity.Id, - WorkflowInstanceId = context.Id, - ActivityType = x.Activity.Type, - ActivityName = x.Activity.Name, - ActivityState = x.ActivityState, - Exception = ExceptionState.FromException(x.Exception), - ActivityTypeVersion = x.Activity.Version, - StartedAt = x.StartedAt, - HasBookmarks = x.Bookmarks.Any(), - Status = x.Status, - CompletedAt = x.CompletedAt + // Get any outcomes that were added to the activity execution context. + var outcomes = activityExecutionContext.JournalData.TryGetValue("Result", out var resultValue) ? resultValue as Outcomes : default; + var payload = new Dictionary(); + + if(outcomes != null) + payload.Add("Outcomes", outcomes); + + // Get any outputs that were added to the activity execution context. + var activity = activityExecutionContext.Activity; + var expressionExecutionContext = activityExecutionContext.ExpressionExecutionContext; + var activityDescriptor = activityExecutionContext.ActivityDescriptor; + var outputDescriptors = activityDescriptor.Outputs; + var outputs = outputDescriptors.ToDictionary(x => x.Name, x => activity.GetOutput(expressionExecutionContext, x.Name)!); + + return new ActivityExecutionRecord + { + Id = activityExecutionContext.Id, + ActivityId = activityExecutionContext.Activity.Id, + WorkflowInstanceId = context.Id, + ActivityType = activityExecutionContext.Activity.Type, + ActivityName = activityExecutionContext.Activity.Name, + ActivityState = activityExecutionContext.ActivityState, + Outputs = outputs, + Payload = payload, + Exception = ExceptionState.FromException(activityExecutionContext.Exception), + ActivityTypeVersion = activityExecutionContext.Activity.Version, + StartedAt = activityExecutionContext.StartedAt, + HasBookmarks = activityExecutionContext.Bookmarks.Any(), + Status = activityExecutionContext.Status, + CompletedAt = activityExecutionContext.CompletedAt + }; }).ToList(); await _activityExecutionStore.SaveManyAsync(entries, context.CancellationToken); diff --git a/update-migrations-runtime.sh b/update-migrations.sh similarity index 100% rename from update-migrations-runtime.sh rename to update-migrations.sh