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); return await strategy.GetPersistenceModeAsync(strategyContext); } 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; } }