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
This commit is contained in:
Sipke Schoorstra 2023-08-23 13:34:30 +02:00 committed by GitHub
parent 1e7f954311
commit 73482a8e4c
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
19 changed files with 170 additions and 42 deletions

View file

@ -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

View file

@ -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
{
/// <inheritdoc />
@ -125,6 +125,12 @@ namespace Elsa.EntityFrameworkCore.MySql.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("longtext");
b.Property<string>("SerializedOutputs")
.HasColumnType("longtext");
b.Property<string>("SerializedPayload")
.HasColumnType("longtext");
b.Property<DateTimeOffset>("StartedAt")
.HasColumnType("datetime(6)");

View file

@ -40,6 +40,10 @@ namespace Elsa.EntityFrameworkCore.MySql.Migrations.Runtime
SerializedActivityState = table.Column<string>(type: "longtext", nullable: true)
.Annotation("MySql:CharSet", "utf8mb4"),
SerializedException = table.Column<string>(type: "longtext", nullable: true)
.Annotation("MySql:CharSet", "utf8mb4"),
SerializedOutputs = table.Column<string>(type: "longtext", nullable: true)
.Annotation("MySql:CharSet", "utf8mb4"),
SerializedPayload = table.Column<string>(type: "longtext", nullable: true)
.Annotation("MySql:CharSet", "utf8mb4")
},
constraints: table =>

View file

@ -122,6 +122,12 @@ namespace Elsa.EntityFrameworkCore.MySql.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("longtext");
b.Property<string>("SerializedOutputs")
.HasColumnType("longtext");
b.Property<string>("SerializedPayload")
.HasColumnType("longtext");
b.Property<DateTimeOffset>("StartedAt")
.HasColumnType("datetime(6)");

View file

@ -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
{
/// <inheritdoc />
@ -128,6 +128,12 @@ namespace Elsa.EntityFrameworkCore.PostgreSql.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("text");
b.Property<string>("SerializedOutputs")
.HasColumnType("text");
b.Property<string>("SerializedPayload")
.HasColumnType("text");
b.Property<DateTimeOffset>("StartedAt")
.HasColumnType("timestamp with time zone");

View file

@ -30,7 +30,9 @@ namespace Elsa.EntityFrameworkCore.PostgreSql.Migrations.Runtime
Status = table.Column<int>(type: "integer", nullable: false),
CompletedAt = table.Column<DateTimeOffset>(type: "timestamp with time zone", nullable: true),
SerializedActivityState = table.Column<string>(type: "text", nullable: true),
SerializedException = table.Column<string>(type: "text", nullable: true)
SerializedException = table.Column<string>(type: "text", nullable: true),
SerializedOutputs = table.Column<string>(type: "text", nullable: true),
SerializedPayload = table.Column<string>(type: "text", nullable: true)
},
constraints: table =>
{

View file

@ -125,6 +125,12 @@ namespace Elsa.EntityFrameworkCore.PostgreSql.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("text");
b.Property<string>("SerializedOutputs")
.HasColumnType("text");
b.Property<string>("SerializedPayload")
.HasColumnType("text");
b.Property<DateTimeOffset>("StartedAt")
.HasColumnType("timestamp with time zone");

View file

@ -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
{
/// <inheritdoc />
@ -128,6 +128,12 @@ namespace Elsa.EntityFrameworkCore.SqlServer.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("nvarchar(max)");
b.Property<string>("SerializedOutputs")
.HasColumnType("nvarchar(max)");
b.Property<string>("SerializedPayload")
.HasColumnType("nvarchar(max)");
b.Property<DateTimeOffset>("StartedAt")
.HasColumnType("datetimeoffset");

View file

@ -30,7 +30,9 @@ namespace Elsa.EntityFrameworkCore.SqlServer.Migrations.Runtime
Status = table.Column<int>(type: "int", nullable: false),
CompletedAt = table.Column<DateTimeOffset>(type: "datetimeoffset", nullable: true),
SerializedActivityState = table.Column<string>(type: "nvarchar(max)", nullable: true),
SerializedException = table.Column<string>(type: "nvarchar(max)", nullable: true)
SerializedException = table.Column<string>(type: "nvarchar(max)", nullable: true),
SerializedOutputs = table.Column<string>(type: "nvarchar(max)", nullable: true),
SerializedPayload = table.Column<string>(type: "nvarchar(max)", nullable: true)
},
constraints: table =>
{

View file

@ -125,6 +125,12 @@ namespace Elsa.EntityFrameworkCore.SqlServer.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("nvarchar(max)");
b.Property<string>("SerializedOutputs")
.HasColumnType("nvarchar(max)");
b.Property<string>("SerializedPayload")
.HasColumnType("nvarchar(max)");
b.Property<DateTimeOffset>("StartedAt")
.HasColumnType("datetimeoffset");

View file

@ -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
{
/// <inheritdoc />
@ -123,6 +123,12 @@ namespace Elsa.EntityFrameworkCore.Sqlite.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("TEXT");
b.Property<string>("SerializedOutputs")
.HasColumnType("TEXT");
b.Property<string>("SerializedPayload")
.HasColumnType("TEXT");
b.Property<string>("StartedAt")
.IsRequired()
.HasColumnType("TEXT");

View file

@ -25,7 +25,9 @@ namespace Elsa.EntityFrameworkCore.Sqlite.Migrations.Runtime
Status = table.Column<int>(type: "INTEGER", nullable: false),
CompletedAt = table.Column<string>(type: "TEXT", nullable: true),
SerializedActivityState = table.Column<string>(type: "TEXT", nullable: true),
SerializedException = table.Column<string>(type: "TEXT", nullable: true)
SerializedException = table.Column<string>(type: "TEXT", nullable: true),
SerializedOutputs = table.Column<string>(type: "TEXT", nullable: true),
SerializedPayload = table.Column<string>(type: "TEXT", nullable: true)
},
constraints: table =>
{

View file

@ -120,6 +120,12 @@ namespace Elsa.EntityFrameworkCore.Sqlite.Migrations.Runtime
b.Property<string>("SerializedException")
.HasColumnType("TEXT");
b.Property<string>("SerializedOutputs")
.HasColumnType("TEXT");
b.Property<string>("SerializedPayload")
.HasColumnType("TEXT");
b.Property<string>("StartedAt")
.IsRequired()
.HasColumnType("TEXT");

View file

@ -38,46 +38,57 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore
/// <inheritdoc />
public async Task<IEnumerable<ActivityExecutionRecord>> FindManyAsync<TOrderBy>(ActivityExecutionRecordFilter filter, ActivityExecutionRecordOrder<TOrderBy> 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();
/// <inheritdoc />
public async Task<IEnumerable<ActivityExecutionRecord>> 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();
/// <inheritdoc />
public async Task<long> CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.CountAsync(queryable => Filter(queryable, filter), cancellationToken);
public async Task<long> CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.CountAsync(queryable => EFCoreActivityExecutionStore.Filter(queryable, filter), cancellationToken);
/// <inheritdoc />
public async Task<long> DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.DeleteWhereAsync(queryable => Filter(queryable, filter), cancellationToken);
public async Task<long> 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<IDictionary<string, object>>(dbContext, entity, "SerializedOutputs");
entity.Exception = DeserializePayload<ExceptionState>(dbContext, entity, "SerializedException");
entity.Payload = DeserializePayload<IDictionary<string, object>>(dbContext, entity, "SerializedPayload");
return ValueTask.CompletedTask;
}
private ValueTask<IDictionary<string, object>?> LoadActivityState(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity)
private IDictionary<string, object>? DeserializeActivityState(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity)
{
var json = dbContext.Entry(entity).Property<string>("SerializedActivityState").CurrentValue;
var dictionary = !string.IsNullOrEmpty(json) ? JsonSerializer.Deserialize<Dictionary<string, JsonElement>>(json) : default;
return ValueTask.FromResult<IDictionary<string, object>?>(dictionary?.ToDictionary(x => x.Key, x => (object)x.Value));
var dictionary = Deserialize<Dictionary<string, JsonElement>>(dbContext, entity, "SerializedActivityState");
return dictionary?.ToDictionary(x => x.Key, x => (object)x.Value);
}
private ValueTask<ExceptionState?> LoadException(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity)
private T? Deserialize<T>(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, string propertyName)
{
var json = dbContext.Entry(entity).Property<string>("SerializedException").CurrentValue;
var exceptionState = !string.IsNullOrEmpty(json) ? _payloadSerializer.Deserialize<ExceptionState>(json) : default;
return ValueTask.FromResult(exceptionState);
var json = dbContext.Entry(entity).Property<string>(propertyName).CurrentValue;
var value = !string.IsNullOrEmpty(json) ? JsonSerializer.Deserialize<T>(json) : default;
return value;
}
private IQueryable<ActivityExecutionRecord> Filter(IQueryable<ActivityExecutionRecord> queryable, ActivityExecutionRecordFilter filter) => filter.Apply(queryable);
private T? DeserializePayload<T>(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, string propertyName)
{
var json = dbContext.Entry(entity).Property<string>(propertyName).CurrentValue;
var payload = !string.IsNullOrEmpty(json) ? _payloadSerializer.Deserialize<T>(json) : default;
return payload;
}
private static IQueryable<ActivityExecutionRecord> Filter(IQueryable<ActivityExecutionRecord> queryable, ActivityExecutionRecordFilter filter) => filter.Apply(queryable);
}

View file

@ -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<string>("SerializedActivityState");
builder.Property<string>("SerializedException");
builder.Property<string>("SerializedPayload");
builder.Property<string>("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)}");

View file

@ -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<IActivityStateSerializer>();
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);

View file

@ -38,6 +38,16 @@ public class ActivityExecutionRecord : Entity
/// 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>
public IDictionary<string, object>? Outputs { get; set; }
/// <summary>
/// Gets or sets the exception that occurred during the activity execution.

View file

@ -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<string, object>();
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);