using System.Text.Json;
using System.Text.Json.Serialization;
using Elsa.Expressions.Contracts;
using Elsa.Expressions.Models;
using Elsa.Extensions;
using Elsa.Workflows.Activities;
using Elsa.Workflows.LogPersistence;
using Elsa.Workflows.LogPersistence.Strategies;
using Elsa.Workflows.Management.Options;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Serialization.Converters;
using Elsa.Workflows.State;
using Humanizer;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace Elsa.Workflows.Runtime;
///
public class DefaultActivityExecutionMapper : IActivityExecutionMapper
{
private readonly JsonSerializerOptions _logPersistenceConfigSerializerOptions;
private readonly IOptions _options;
private readonly IExpressionEvaluator _expressionEvaluator;
private readonly ILogger _logger;
private readonly IDictionary _logPersistenceStrategies;
public DefaultActivityExecutionMapper(
IOptions options,
ILogPersistenceStrategyService logPersistenceStrategyService,
IExpressionEvaluator expressionEvaluator,
IExpressionDescriptorRegistry expressionDescriptorRegistry,
ILogger logger)
{
_options = options;
_expressionEvaluator = expressionEvaluator;
_logger = logger;
_logPersistenceStrategies = logPersistenceStrategyService.ListStrategies().ToDictionary(x => x.GetType().GetSimpleAssemblyQualifiedName(), x => x);
_logPersistenceConfigSerializerOptions = new JsonSerializerOptions
{
PropertyNameCaseInsensitive = true
}.WithConverters(
new ExpressionJsonConverterFactory(expressionDescriptorRegistry),
new JsonStringEnumConverter(),
new ExpandoObjectConverterFactory());
}
private const string LegacyLogPersistenceModeKey = "logPersistenceMode";
private const string LogPersistenceConfigKey = "logPersistenceConfig";
///
public async Task MapAsync(ActivityExecutionContext source)
{
/* The following legacy JSON structure is expected to be found in the custom properties of the workflow and activity:
* {
* "logPersistenceMode": {
* "default": "default",
* "inputs": { k : v },
* "outputs": { k: v }
* }
* }
*/
/* The following JSON structure is expected to be found in the custom properties of the workflow and activity:
* {
* "logPersistenceConfig": {
* "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": "..." } }
* }
* }
*/
var cancellationToken = source.WorkflowExecutionContext.CancellationToken;
var workflow = (Workflow?)source.GetAncestors().FirstOrDefault(x => x.Activity is Workflow)?.Activity ?? source.WorkflowExecutionContext.Workflow;
var rootActivityExecutionContext = source.WorkflowExecutionContext.ActivityExecutionContexts.First(x => x.ParentActivityExecutionContext == null);
var workflowPersistenceProperty = await GetDefaultPersistenceModeAsync(rootActivityExecutionContext.ExpressionExecutionContext, workflow.CustomProperties, () => _options.Value.LogPersistenceMode, cancellationToken);
var activityPersistencePropertyDefault = await GetDefaultPersistenceModeAsync(source.ExpressionExecutionContext, source.Activity.CustomProperties, () => workflowPersistenceProperty, cancellationToken);
var legacyActivityPersistenceProperties = source.Activity.CustomProperties.GetValueOrDefault>(LegacyLogPersistenceModeKey, () => new Dictionary());
var activityPersistenceProperties = source.Activity.CustomProperties.GetValueOrDefault>(LogPersistenceConfigKey, () => new Dictionary());
var payload = GetPayload(source);
var outputs = GetOutputs(source);
outputs = await StorePropertyUsingPersistenceMode(
source.ExpressionExecutionContext,
outputs,
legacyActivityPersistenceProperties!.GetValueOrDefault("outputs", () => new Dictionary())!,
activityPersistenceProperties!.GetValueOrDefault("outputs", () => new Dictionary())!,
activityPersistencePropertyDefault,
cancellationToken);
var inputs = await StorePropertyUsingPersistenceMode(
source.ExpressionExecutionContext,
source.ActivityState,
legacyActivityPersistenceProperties!.GetValueOrDefault("inputs", () => new Dictionary())!,
activityPersistenceProperties!.GetValueOrDefault("inputs", () => new Dictionary())!,
activityPersistencePropertyDefault,
cancellationToken);
return new ActivityExecutionRecord
{
Id = source.Id,
ActivityId = source.Activity.Id,
ActivityNodeId = source.NodeId,
WorkflowInstanceId = source.WorkflowExecutionContext.Id,
ActivityType = source.Activity.Type,
ActivityName = source.Activity.Name,
ActivityState = inputs,
Outputs = outputs,
Properties = source.Properties,
Payload = payload,
Exception = ExceptionState.FromException(source.Exception),
ActivityTypeVersion = source.Activity.Version,
StartedAt = source.StartedAt,
HasBookmarks = source.Bookmarks.Any(),
Status = GetAggregateStatus(source),
CompletedAt = source.CompletedAt
};
}
private async Task GetDefaultPersistenceModeAsync(ExpressionExecutionContext expressionExecutionContext, IDictionary customProperties, Func defaultFactory, CancellationToken cancellationToken)
{
var legacyProperties = customProperties.GetValueOrDefault>(LegacyLogPersistenceModeKey, () => new Dictionary());
var properties = customProperties.GetValueOrDefault>(LogPersistenceConfigKey, () => new Dictionary());
var defaultPersistenceConfigObject = properties?.TryGetValue("default", out var defaultPersistenceConfigObjectValue) == true ? defaultPersistenceConfigObjectValue : null;
if (defaultPersistenceConfigObject == null)
{
var legacyPersistencePropertyDefault = legacyProperties!.GetValueOrDefault("default", defaultFactory);
if (legacyPersistencePropertyDefault == LogPersistenceMode.Inherit)
return defaultFactory();
return legacyPersistencePropertyDefault;
}
var defaultPersistenceConfig = Convert(defaultPersistenceConfigObject);
return await EvaluateLogPersistenceConfigAsync(defaultPersistenceConfig, expressionExecutionContext, defaultFactory, cancellationToken);
}
private async Task> StorePropertyUsingPersistenceMode(
ExpressionExecutionContext expressionExecutionContext,
IDictionary state,
IDictionary obsoletePersistenceModeConfiguration,
IDictionary persistenceStrategyConfiguration,
LogPersistenceMode defaultLogPersistenceMode,
CancellationToken cancellationToken)
{
var result = new Dictionary();
foreach (var value in state)
{
var propKey = value.Key.Camelize();
var logPersistenceConfigObject = persistenceStrategyConfiguration.GetValueOrDefault(propKey, () => null);
var logPersistenceConfig = Convert(logPersistenceConfigObject);
var mode = defaultLogPersistenceMode;
if (logPersistenceConfig != null)
{
mode = await EvaluateLogPersistenceConfigAsync(logPersistenceConfig, expressionExecutionContext, () => defaultLogPersistenceMode, cancellationToken);
}
else
{
mode = obsoletePersistenceModeConfiguration.GetValueOrDefault(propKey, () => defaultLogPersistenceMode);
}
if (mode == LogPersistenceMode.Include || mode == LogPersistenceMode.Inherit && defaultLogPersistenceMode is LogPersistenceMode.Include or LogPersistenceMode.Inherit)
result.Add(value.Key, value.Value);
}
return result;
}
private LogPersistenceConfiguration? Convert(object? value)
{
if (value == null)
return null;
if (value is LogPersistenceConfiguration c)
return c;
var json = JsonSerializer.Serialize(value);
var config = JsonSerializer.Deserialize(json, _logPersistenceConfigSerializerOptions);
return config;
}
private async Task EvaluateLogPersistenceConfigAsync(LogPersistenceConfiguration? config, ExpressionExecutionContext executionContext, Func defaultMode, CancellationToken cancellationToken)
{
if (config == null)
return defaultMode();
if (config.EvaluationMode == LogPersistenceEvaluationMode.Strategy)
{
var strategyTypeName = config.StrategyType ?? typeof(Inherit).GetSimpleAssemblyQualifiedName();
var strategy = _logPersistenceStrategies.TryGetValue(strategyTypeName, out var v) ? v : null;
if (strategy == null)
return defaultMode();
var strategyContext = new LogPersistenceStrategyContext(cancellationToken);
var logMode = await strategy.GetPersistenceModeAsync(strategyContext);
return logMode == LogPersistenceMode.Inherit ? defaultMode() : logMode;
}
if (config.Expression == null)
return defaultMode();
var expression = config.Expression;
try
{
return await _expressionEvaluator.EvaluateAsync(expression, executionContext);
}
catch (Exception e)
{
_logger.LogWarning(e, "Error evaluating log persistence expression");
return defaultMode();
}
}
private static ActivityStatus GetAggregateStatus(ActivityExecutionContext context)
{
// If any child activity is faulted, the aggregate status is faulted.
var descendantContexts = context.GetDescendants().ToList();
if (descendantContexts.Any(x => x.Status == ActivityStatus.Faulted))
return ActivityStatus.Faulted;
return context.Status;
}
private static IDictionary GetPayload(ActivityExecutionContext source)
{
var outcomes = source.JournalData.TryGetValue("Outcomes", out var resultValue) ? resultValue as string[] : default;
var payload = new Dictionary();
if (outcomes != null)
payload.Add("Outcomes", outcomes);
return payload;
}
private static IDictionary GetOutputs(ActivityExecutionContext source)
{
var activity = source.Activity;
var expressionExecutionContext = source.ExpressionExecutionContext;
var activityDescriptor = source.ActivityDescriptor;
var outputDescriptors = activityDescriptor.Outputs;
var outputs = outputDescriptors.ToDictionary(x => x.Name, x =>
{
if (x.IsSerializable == false)
return "(not serializable)";
var cachedValue = activity.GetOutput(expressionExecutionContext, x.Name);
if (cachedValue != default)
return cachedValue;
if (x.ValueGetter(activity) is Output output && source.TryGet(output.MemoryBlockReference(), out var outputValue))
return outputValue;
return default;
});
return outputs;
}
}