Implement internal state activity persistence and logging mechanisms (#6601)

* Implement internal state activity persistence and logging mechanisms

Updated property handling to support nullable dictionaries and improved persistable states. Adjusted serialization logic to handle optional fields more robustly, ensuring better compatibility with log persistence mappings and internal state evaluations.

* Replace default! with null! for string properties

Updated string properties in various records to use null! instead of default! for consistency and clarity. Additionally, adjusted methods to check collection existence before serialization and streamlined object initializations with simplified syntax where possible.

* Fix nullable types in DeserializeActivityState method

Updated the method's return type and JSON deserialization to properly handle nullable values. This ensures better alignment with the method's behavior and avoids potential null reference issues.

* Fix null reference issues in InputOutputLoggingTests

Replaced forced null dereferences with safe navigation checks to prevent potential null reference exceptions. This ensures more robust and error-free test execution for activity state validations.
This commit is contained in:
Sipke Schoorstra 2025-04-18 14:29:18 +02:00 committed by GitHub
parent 67f3ceb801
commit 476656ccce
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 83 additions and 63 deletions

View file

@ -50,14 +50,13 @@ public class ConnectionMiddleware(ActivityMiddlewareDelegate next
LogConnectionExtensions.LogConnectionIsNull(logger);
else
{
//Get connection from store, if exist,
var connectionConfiguration = await connectionStore.FindAsync(new Persistence.Filters.ConnectionDefinitionFilter() { Name = connectionName });
// Get connection from store, if exists.
var connectionConfiguration = await connectionStore.FindAsync(new() { Name = connectionName });
if (connectionConfiguration != null)
{
dynamic deserializedjson = JsonSerializer.Deserialize(connectionConfiguration.ConnectionConfiguration, propertyType, SerializerOptions);
inputValue.Properties = deserializedjson;
dynamic deserializedJson = connectionConfiguration.ConnectionConfiguration.Deserialize(propertyType, SerializerOptions)!;
inputValue.Properties = deserializedJson;
input.ValueSetter(context.Activity, inputValue);
}
else

View file

@ -10,22 +10,22 @@ internal class ActivityExecutionRecordRecord : Record
/// <summary>
/// Gets or sets the workflow instance ID.
/// </summary>
public string WorkflowInstanceId { get; set; } = default!;
public string WorkflowInstanceId { get; set; } = null!;
/// <summary>
/// Gets or sets the activity ID.
/// </summary>
public string ActivityId { get; set; } = default!;
public string ActivityId { get; set; } = null!;
/// <summary>
/// Gets or sets the activity node ID.
/// </summary>
public string ActivityNodeId { get; set; } = default!;
public string ActivityNodeId { get; set; } = null!;
/// <summary>
/// The type of the activity.
/// </summary>
public string ActivityType { get; set; } = default!;
public string ActivityType { get; set; } = null!;
/// <summary>
/// The version of the activity type.
@ -75,7 +75,7 @@ internal class ActivityExecutionRecordRecord : Record
/// <summary>
/// Gets or sets the status of the activity.
/// </summary>
public string Status { get; set; } = default!;
public string Status { get; set; } = null!;
/// <summary>
/// Gets or sets the time at which the activity execution completed.

View file

@ -10,22 +10,22 @@ internal class ActivityExecutionSummaryRecord : Record
/// <summary>
/// Gets or sets the workflow instance ID.
/// </summary>
public string WorkflowInstanceId { get; set; } = default!;
public string WorkflowInstanceId { get; set; } = null!;
/// <summary>
/// Gets or sets the activity ID.
/// </summary>
public string ActivityId { get; set; } = default!;
public string ActivityId { get; set; } = null!;
/// <summary>
/// Gets or sets the activity node ID.
/// </summary>
public string ActivityNodeId { get; set; } = default!;
public string ActivityNodeId { get; set; } = null!;
/// <summary>
/// The type of the activity.
/// </summary>
public string ActivityType { get; set; } = default!;
public string ActivityType { get; set; } = null!;
/// <summary>
/// The version of the activity type.
@ -50,7 +50,7 @@ internal class ActivityExecutionSummaryRecord : Record
/// <summary>
/// Gets or sets the status of the activity.
/// </summary>
public string Status { get; set; } = default!;
public string Status { get; set; } = null!;
/// <summary>
/// Gets or sets the time at which the activity execution completed.

View file

@ -4,5 +4,5 @@ namespace Elsa.Dapper.Modules.Runtime.Records;
internal class KeyValuePairRecord : Record
{
public string Value { get; set; } = default!;
public string Value { get; set; } = null!;
}

View file

@ -4,9 +4,9 @@ namespace Elsa.Dapper.Modules.Runtime.Records;
internal class StoredBookmarkRecord : Record
{
public string ActivityTypeName { get; set; } = default!;
public string Hash { get; set; } = default!;
public string WorkflowInstanceId { get; set; } = default!;
public string ActivityTypeName { get; set; } = null!;
public string Hash { get; set; } = null!;
public string WorkflowInstanceId { get; set; } = null!;
public string? CorrelationId { get; set; }
public string? ActivityInstanceId { get; set; }
public string? SerializedPayload { get; set; }

View file

@ -4,10 +4,10 @@ namespace Elsa.Dapper.Modules.Runtime.Records;
internal class StoredTriggerRecord : Record
{
public string WorkflowDefinitionId { get; set; } = default!;
public string WorkflowDefinitionVersionId { get; set; } = default!;
public string Name { get; set; } = default!;
public string ActivityId { get; set; } = default!;
public string WorkflowDefinitionId { get; set; } = null!;
public string WorkflowDefinitionVersionId { get; set; } = null!;
public string Name { get; set; } = null!;
public string ActivityId { get; set; } = null!;
public string? Hash { get; set; }
public string? SerializedPayload { get; set; }
}

View file

@ -4,18 +4,18 @@ namespace Elsa.Dapper.Modules.Runtime.Records;
internal class WorkflowExecutionLogRecordRecord : Record
{
public string Id { get; set; } = default!;
public string WorkflowDefinitionId { get; set; } = default!;
public string WorkflowDefinitionVersionId { get; set; } = default!;
public string WorkflowInstanceId { get; set; } = default!;
public string Id { get; set; } = null!;
public string WorkflowDefinitionId { get; set; } = null!;
public string WorkflowDefinitionVersionId { get; set; } = null!;
public string WorkflowInstanceId { get; set; } = null!;
public int WorkflowVersion { get; set; }
public string ActivityInstanceId { get; set; } = default!;
public string ActivityInstanceId { get; set; } = null!;
public string? ParentActivityInstanceId { get; set; }
public string ActivityId { get; set; } = default!;
public string ActivityType { get; set; } = default!;
public string ActivityId { get; set; } = null!;
public string ActivityType { get; set; } = null!;
public int ActivityTypeVersion { get; set; }
public string? ActivityName { get; set; } = default!;
public string ActivityNodeId { get; set; } = default!;
public string? ActivityName { get; set; }
public string ActivityNodeId { get; set; } = null!;
public DateTimeOffset Timestamp { get; set; }
public long Sequence { get; set; }
public string? EventName { get; set; }

View file

@ -109,7 +109,7 @@ internal class DapperActivityExecutionRecordStore(Store<ActivityExecutionRecordR
private ActivityExecutionRecordRecord Map(ActivityExecutionRecord source)
{
return new ActivityExecutionRecordRecord
return new()
{
Id = source.Id,
ActivityId = source.ActivityId,
@ -122,18 +122,18 @@ internal class DapperActivityExecutionRecordStore(Store<ActivityExecutionRecordR
HasBookmarks = source.HasBookmarks,
Status = source.Status.ToString(),
ActivityTypeVersion = source.ActivityTypeVersion,
SerializedActivityState = source.ActivityState != null ? safeSerializer.Serialize(source.ActivityState) : null,
SerializedPayload = source.Payload != null ? safeSerializer.Serialize(source.Payload) : null,
SerializedOutputs = source.Outputs != null ? safeSerializer.Serialize(source.Outputs) : null,
SerializedActivityState = source.ActivityState?.Any() == true ? safeSerializer.Serialize(source.ActivityState) : null,
SerializedPayload = source.Payload?.Any() == true ? safeSerializer.Serialize(source.Payload) : null,
SerializedOutputs = source.Outputs?.Any() == true ? safeSerializer.Serialize(source.Outputs) : null,
SerializedException = source.Exception != null ? payloadSerializer.Serialize(source.Exception) : null,
SerializedProperties = source.Properties.Any() ? safeSerializer.Serialize(source.Properties) : null,
SerializedProperties = source.Properties?.Any() == true ? safeSerializer.Serialize(source.Properties) : null,
TenantId = source.TenantId
};
}
private ActivityExecutionRecord Map(ActivityExecutionRecordRecord source)
{
return new ActivityExecutionRecord
return new()
{
Id = source.Id,
ActivityId = source.ActivityId,
@ -146,18 +146,18 @@ internal class DapperActivityExecutionRecordStore(Store<ActivityExecutionRecordR
HasBookmarks = source.HasBookmarks,
Status = Enum.Parse<ActivityStatus>(source.Status),
ActivityTypeVersion = source.ActivityTypeVersion,
ActivityState = source.SerializedActivityState != null ? payloadSerializer.Deserialize<IDictionary<string, object>>(source.SerializedActivityState) : default,
Payload = source.SerializedPayload != null ? safeSerializer.Deserialize<IDictionary<string, object>>(source.SerializedPayload) : default,
Outputs = source.SerializedOutputs != null ? safeSerializer.Deserialize<IDictionary<string, object?>>(source.SerializedOutputs) : default,
Exception = source.SerializedException != null ? payloadSerializer.Deserialize<ExceptionState>(source.SerializedException) : default,
Properties = source.SerializedProperties != null ? safeSerializer.Deserialize<IDictionary<string, object>>(source.SerializedProperties) : new Dictionary<string, object>(),
ActivityState = source.SerializedActivityState != null ? payloadSerializer.Deserialize<IDictionary<string, object?>>(source.SerializedActivityState) : null,
Payload = source.SerializedPayload != null ? safeSerializer.Deserialize<IDictionary<string, object>>(source.SerializedPayload) : null,
Outputs = source.SerializedOutputs != null ? safeSerializer.Deserialize<IDictionary<string, object?>>(source.SerializedOutputs) : null,
Exception = source.SerializedException != null ? payloadSerializer.Deserialize<ExceptionState>(source.SerializedException) : null,
Properties = source.SerializedProperties != null ? safeSerializer.Deserialize<IDictionary<string, object>>(source.SerializedProperties) : null,
TenantId = source.TenantId
};
}
private ActivityExecutionRecordSummary MapSummary(ActivityExecutionSummaryRecord source)
{
return new ActivityExecutionRecordSummary
return new()
{
Id = source.Id,
ActivityId = source.ActivityId,

View file

@ -85,13 +85,13 @@ public class EFCoreActivityExecutionStore(
{
entity = entity.SanitizeLogMessage();
var compressionAlgorithm = options.Value.CompressionAlgorithm ?? nameof(None);
var serializedActivityState = entity.ActivityState != null ? safeSerializer.Serialize(entity.ActivityState) : null;
var serializedActivityState = entity.ActivityState?.Count > 0 ? safeSerializer.Serialize(entity.ActivityState) : null;
var compressedSerializedActivityState = serializedActivityState != null ? await compressionCodecResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : 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 = entity.Properties.Any() ? payloadSerializer.Serialize(entity.Properties) : null;
dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = entity.Properties?.Any() == true ? payloadSerializer.Serialize(entity.Properties) : null;
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;
}
@ -104,13 +104,13 @@ public class EFCoreActivityExecutionStore(
entity.ActivityState = await DeserializeActivityState(dbContext, entity, cancellationToken);
entity.Outputs = Deserialize<IDictionary<string, object?>>(dbContext, entity, "SerializedOutputs");
entity.Properties = DeserializePayload<IDictionary<string, object>?>(dbContext, entity, "SerializedProperties") ?? new Dictionary<string, object>();
entity.Properties = DeserializePayload<IDictionary<string, object>?>(dbContext, entity, "SerializedProperties");
entity.Exception = DeserializePayload<ExceptionState>(dbContext, entity, "SerializedException");
entity.Payload = DeserializePayload<IDictionary<string, object>>(dbContext, entity, "SerializedPayload");
}
[RequiresUnreferencedCode("Calls System.Text.Json.JsonSerializer.Deserialize<TValue>(String, JsonSerializerOptions)")]
private async Task<IDictionary<string, object>?> DeserializeActivityState(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken)
private async Task<IDictionary<string, object?>?> DeserializeActivityState(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken)
{
var json = dbContext.Entry(entity).Property<string>("SerializedActivityState").CurrentValue;
@ -119,11 +119,11 @@ 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);
return dictionary?.ToDictionary(x => x.Key, x => (object)x.Value);
var dictionary = JsonSerializer.Deserialize<IDictionary<string, object?>>(json);
return dictionary?.ToDictionary(x => x.Key, x => x.Value);
}
return default;
return null;
}
[RequiresUnreferencedCode("Calls System.Text.Json.JsonSerializer.Deserialize<TValue>(String, JsonSerializerOptions)")]

View file

@ -57,7 +57,7 @@ public class ActivityExecutionRecord : Entity, ILogRecord
/// <summary>
/// Any properties provided by the activity.
/// </summary>
public IDictionary<string, object> Properties { get; set; } = new Dictionary<string, object>();
public IDictionary<string, object>? Properties { get; set; }
/// <summary>
/// Gets or sets the exception that occurred during the activity execution.

View file

@ -6,4 +6,5 @@ public class ActivityLogPersistenceModeMap
{
public IDictionary<string, LogPersistenceMode> Inputs { get; set; } = new Dictionary<string, LogPersistenceMode>();
public IDictionary<string, LogPersistenceMode> Outputs { get; set; } = new Dictionary<string, LogPersistenceMode>();
public LogPersistenceMode InternalState { get; set; }
}

View file

@ -31,6 +31,7 @@ namespace Elsa.Workflows.Runtime;
* "default": { "evaluationMode": "Strategy", "strategyType": "Elsa.Workflows.LogPersistence.Strategies.Inherit, Elsa.Workflows.Core", "expression": "..." },
* "inputs": { "input1" : { "evaluationMode": "Strategy", "strategyType": "Elsa.Workflows.LogPersistence.Strategies.Inherit, Elsa.Workflows.Core", "expression": "..." } },
* "outputs": { "output1" : { "evaluationMode": "Strategy", "strategyType": "Elsa.Workflows.LogPersistence.Strategies.Inherit, Elsa.Workflows.Core", "expression": "..." } }
* "internalState": { "evaluationMode": "Strategy", "strategyType": "Elsa.Workflows.LogPersistence.Strategies.Inherit, Elsa.Workflows.Core", "expression": "..." }
* }
* }
*/
@ -73,6 +74,7 @@ public class ActivityPropertyLogPersistenceEvaluator : IActivityPropertyLogPersi
await EvaluatePropertiesAsync(context, "inputs", context.ActivityDescriptor.Inputs, legacyProps, configProps, defaultMode, map.Inputs, cancellationToken);
await EvaluatePropertiesAsync(context, "outputs", context.ActivityDescriptor.Outputs, legacyProps, configProps, defaultMode, map.Outputs, cancellationToken);
map.InternalState = await EvaluateInternalStateModeAsync(context.ExpressionExecutionContext, context.Activity.CustomProperties, defaultMode, cancellationToken);
return map;
}
@ -85,8 +87,7 @@ public class ActivityPropertyLogPersistenceEvaluator : IActivityPropertyLogPersi
return await GetPersistablePropertiesAsync(context, outputs, "outputs", legacyProps, configProps, defaultMode, cancellationToken);
}
private async Task<(IDictionary<string, object> legacyProps, IDictionary<string, object> configProps, LogPersistenceMode defaultMode)>
GetPersistenceDefaultsAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
private async Task<(IDictionary<string, object> legacyProps, IDictionary<string, object> configProps, LogPersistenceMode defaultMode)> GetPersistenceDefaultsAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
{
var legacyProps = context.Activity.CustomProperties.GetValueOrDefault<IDictionary<string, object>>(LegacyKey, () => new Dictionary<string, object>())!;
var rootContext = context.WorkflowExecutionContext.ActivityExecutionContexts.First(x => x.ParentActivityExecutionContext == null);
@ -133,6 +134,18 @@ public class ActivityPropertyLogPersistenceEvaluator : IActivityPropertyLogPersi
var mode = legacySection.GetValueOrDefault(key, () => defaultMode);
return ResolveMode(mode, () => defaultMode);
}
private async Task<LogPersistenceMode> EvaluateInternalStateModeAsync(
ExpressionExecutionContext executionContext,
IDictionary<string, object> currentConfig,
LogPersistenceMode defaultMode,
CancellationToken cancellationToken)
{
var configObject = currentConfig.GetValueOrDefault("internalState", () => new Dictionary<string, object>())!;
var config = ConvertToConfig(configObject);
if (config != null) return await EvaluateConfigAsync(config, executionContext, () => defaultMode, cancellationToken);
return LogPersistenceMode.Inherit;
}
private async Task<Dictionary<string, object>> GetPersistablePropertiesAsync(
ActivityExecutionContext context,

View file

@ -13,8 +13,10 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper
var outputs = source.GetOutputs();
var inputs = source.GetInputs();
var persistenceMap = source.GetLogPersistenceModeMap();
var persistableInputs = GetPersistableProperties(inputs, persistenceMap.Inputs);
var persistableOutputs = GetPersistableProperties(outputs, persistenceMap.Outputs);
var persistableInputs = GetPersistableInputOutput(inputs, persistenceMap.Inputs);
var persistableOutputs = GetPersistableInputOutput(outputs, persistenceMap.Outputs);
var persistableProperties = GetPersistableDictionary(source.Properties!, persistenceMap.InternalState);
var persistablePayload = GetPersistableDictionary(payload!, persistenceMap.InternalState);
return new()
{
@ -26,8 +28,8 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper
ActivityName = source.Activity.Name,
ActivityState = persistableInputs,
Outputs = persistableOutputs,
Properties = source.Properties,
Payload = payload,
Properties = persistableProperties,
Payload = persistablePayload!,
Exception = ExceptionState.FromException(source.Exception),
ActivityTypeVersion = source.Activity.Version,
StartedAt = source.StartedAt,
@ -43,7 +45,7 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper
return Task.FromResult(Map(source));
}
private IDictionary<string, object?> GetPersistableProperties(IDictionary<string, object> state, IDictionary<string, LogPersistenceMode> map)
private IDictionary<string, object?> GetPersistableInputOutput(IDictionary<string, object> state, IDictionary<string, LogPersistenceMode> map)
{
var result = new Dictionary<string, object?>();
foreach (var stateEntry in state)
@ -55,7 +57,12 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper
return result;
}
private IDictionary<string, object?>? GetPersistableDictionary(IDictionary<string, object?> dictionary, LogPersistenceMode mode)
{
return mode == LogPersistenceMode.Include ? dictionary : null;
}
private static IDictionary<string, object> GetPayload(ActivityExecutionContext source)
{
var outcomes = source.JournalData.TryGetValue("Outcomes", out var resultValue) ? resultValue as string[] : null;

View file

@ -32,7 +32,7 @@ public class InputOutputLoggingTests(App app) : AppComponentTest(app)
{
var activityExecutionRecord = activityExecutionRecords[i];
var shouldBeIncluded = shouldBeIncludedArray[i];
var isIncluded = activityExecutionRecord.ActivityState!.ContainsKey(nameof(WriteLine.Text));
var isIncluded = activityExecutionRecord.ActivityState?.ContainsKey(nameof(WriteLine.Text)) == true;
Assert.Equal(shouldBeIncluded, isIncluded);
}
}
@ -55,8 +55,8 @@ public class InputOutputLoggingTests(App app) : AppComponentTest(app)
await ExecuteWorkflowAsync("input-output-logging-3");
var setOutput1Record = await GetRecordByActivityNameAsync("SetOutput1");
var setOutput2Record = await GetRecordByActivityNameAsync("SetOutput2");
var output1IsIncluded = setOutput1Record.ActivityState!.ContainsKey("OutputName");
var output2IsIncluded = setOutput2Record.ActivityState!.ContainsKey("OutputName");
var output1IsIncluded = setOutput1Record.ActivityState?.ContainsKey("OutputName") == true;
var output2IsIncluded = setOutput2Record.ActivityState?.ContainsKey("OutputName") == true;
Assert.False(output1IsIncluded);
Assert.True(output2IsIncluded);