< Summary

Information
Class: Elsa.Workflows.WorkflowStateExtractor
Assembly: Elsa.Workflows.Core
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs
Line coverage
96%
Covered lines: 209
Uncovered lines: 8
Coverable lines: 217
Total lines: 348
Line coverage: 96.3%
Branch coverage
75%
Covered branches: 45
Total branches: 60
Branch coverage: 75%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs

#LineLine coverage
 1using Elsa.Extensions;
 2using Elsa.Workflows.Models;
 3using Elsa.Workflows.Services;
 4using Elsa.Workflows.State;
 5using Microsoft.Extensions.Logging;
 6
 7namespace Elsa.Workflows;
 8
 9/// <inheritdoc />
 61310public class WorkflowStateExtractor(ILogger<WorkflowStateExtractor> logger) : IWorkflowStateExtractor
 11{
 12    /// <inheritdoc />
 13    public WorkflowState Extract(WorkflowExecutionContext workflowExecutionContext)
 14    {
 51615        var state = new WorkflowState
 51616        {
 51617            Id = workflowExecutionContext.Id,
 51618            DefinitionId = workflowExecutionContext.Workflow.Identity.DefinitionId,
 51619            DefinitionVersionId = workflowExecutionContext.Workflow.Identity.Id,
 51620            DefinitionVersion = workflowExecutionContext.Workflow.Identity.Version,
 51621            CorrelationId = workflowExecutionContext.CorrelationId,
 51622            Name = workflowExecutionContext.Name,
 51623            ParentWorkflowInstanceId = workflowExecutionContext.ParentWorkflowInstanceId,
 51624            Status = workflowExecutionContext.Status,
 51625            SubStatus = workflowExecutionContext.SubStatus,
 51626            IsExecuting = workflowExecutionContext.IsExecuting,
 51627            Bookmarks = workflowExecutionContext.Bookmarks,
 51628            ExecutionLogSequence = workflowExecutionContext.ExecutionLogSequence,
 51629            Input = GetPersistableInput(workflowExecutionContext),
 51630            Output = workflowExecutionContext.Output,
 51631            Incidents = workflowExecutionContext.Incidents,
 51632            IsSystem = workflowExecutionContext.Workflow.IsSystem,
 51633            CreatedAt = workflowExecutionContext.CreatedAt,
 51634            UpdatedAt = workflowExecutionContext.UpdatedAt,
 51635            FinishedAt = workflowExecutionContext.FinishedAt,
 51636        };
 37
 51638        ExtractProperties(state, workflowExecutionContext);
 51639        ExtractActiveActivityExecutionContexts(state, workflowExecutionContext);
 51640        ExtractCompletionCallbacks(state, workflowExecutionContext);
 51641        ExtractScheduledActivities(state, workflowExecutionContext);
 42
 51643        return state;
 44    }
 45
 46    /// <inheritdoc />
 47    public async Task<WorkflowExecutionContext> ApplyAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowSt
 48    {
 14149        workflowExecutionContext.Id = state.Id;
 14150        workflowExecutionContext.CorrelationId = state.CorrelationId;
 14151        workflowExecutionContext.Name = state.Name;
 14152        workflowExecutionContext.ParentWorkflowInstanceId = state.ParentWorkflowInstanceId;
 14153        workflowExecutionContext.SubStatus = state.SubStatus;
 14154        workflowExecutionContext.IsExecuting = state.IsExecuting;
 14155        workflowExecutionContext.Bookmarks = state.Bookmarks;
 14156        workflowExecutionContext.Output = state.Output;
 14157        workflowExecutionContext.ExecutionLogSequence = state.ExecutionLogSequence;
 14158        workflowExecutionContext.CreatedAt = state.CreatedAt;
 14159        workflowExecutionContext.UpdatedAt = state.UpdatedAt;
 14160        workflowExecutionContext.FinishedAt = state.FinishedAt;
 14161        ApplyInput(state, workflowExecutionContext);
 14162        ApplyProperties(state, workflowExecutionContext);
 14163        await ApplyActivityExecutionContextsAsync(state, workflowExecutionContext);
 14164        ApplyCompletionCallbacks(state, workflowExecutionContext);
 14165        ApplyScheduledActivities(state, workflowExecutionContext);
 14166        return workflowExecutionContext;
 14167    }
 68
 69    private void ApplyInput(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
 70    {
 71        // Only add input from state if the input doesn't already exist on the workflow execution context.
 36072        foreach (var inputItem in state.Input)
 3973            if (!workflowExecutionContext.Input.ContainsKey(inputItem.Key))
 074                workflowExecutionContext.Input.Add(inputItem.Key, inputItem.Value);
 14175    }
 76
 77    private IDictionary<string, object> GetPersistableInput(WorkflowExecutionContext workflowExecutionContext)
 78    {
 79        // TODO: This is a temporary solution. We need to find a better way to handle this.
 53080        var persistableInput = workflowExecutionContext.Workflow.Inputs.Where(x => x.StorageDriverType == typeof(Workflo
 51681        var input = workflowExecutionContext.Input;
 51682        var filteredInput = new Dictionary<string, object>();
 83
 103284        foreach (var inputDefinition in persistableInput)
 85        {
 086            if (input.TryGetValue(inputDefinition.Name, out var value))
 087                filteredInput.Add(inputDefinition.Name, value);
 88        }
 89
 51690        return filteredInput;
 91    }
 92
 93    private void ExtractProperties(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
 94    {
 51695        state.Properties = workflowExecutionContext.Properties;
 51696    }
 97
 98    private void ApplyProperties(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
 99    {
 100        // Merge properties.
 350101        foreach (var property in state.Properties)
 34102            workflowExecutionContext.Properties[property.Key] = property.Value;
 141103    }
 104
 105    private async Task ApplyActivityExecutionContextsAsync(WorkflowState state, WorkflowExecutionContext workflowExecuti
 106    {
 107        var activityExecutionContexts = (await Task.WhenAll(state.ActivityExecutionContexts.Select(async item => await C
 108
 109        var lookup = activityExecutionContexts.ToDictionary(x => x.Id);
 110
 111        // Reconstruct hierarchy.
 112        foreach (var contextState in state.ActivityExecutionContexts.Where(x => !string.IsNullOrWhiteSpace(x.ParentConte
 113        {
 97114            var parentContextId = contextState.ParentContextId;
 97115            if (parentContextId == null || !lookup.TryGetValue(parentContextId, out var parentContext))
 116            {
 0117                logger.LogWarning("Parent context with ID '{ParentContextId}' not found for context with ID '{ContextId}
 0118                continue; // Skip if parent context is not found.
 119            }
 120
 97121            var contextId = contextState.Id;
 122
 97123            if (lookup.TryGetValue(contextId, out var context))
 124            {
 97125                context.ExpressionExecutionContext.ParentContext = parentContext.ExpressionExecutionContext;
 97126                context.ParentActivityExecutionContext = parentContext;
 127            }
 128        }
 129
 130        // Assign root expression execution context.
 131        var rootActivityExecutionContexts = activityExecutionContexts.Where(x => x.ExpressionExecutionContext.ParentCont
 132
 282133        foreach (var rootActivityExecutionContext in rootActivityExecutionContexts)
 0134            rootActivityExecutionContext.ExpressionExecutionContext.ParentContext = workflowExecutionContext.ExpressionE
 135
 141136        workflowExecutionContext.ActivityExecutionContexts = activityExecutionContexts;
 141137        return;
 138
 139        async Task<ActivityExecutionContext?> CreateActivityExecutionContextAsync(ActivityExecutionContextState activity
 140        {
 143141            var activity = workflowExecutionContext.FindActivityByNodeId(activityExecutionContextState.ScheduledActivity
 142
 143            // Activity can be null in case the workflow instance was migrated to a newer version that no longer contain
 143144            if (activity == null)
 145            {
 2146                LogSkippedStateReference(
 2147                    state,
 2148                    workflowExecutionContext,
 2149                    "ActivityExecutionContext",
 2150                    activityExecutionContextState.Id,
 2151                    activityExecutionContextState.ScheduledActivityNodeId);
 2152                return null;
 153            }
 154
 141155            var properties = activityExecutionContextState.Properties;
 141156            var metadata = activityExecutionContextState.Metadata;
 141157            var activityExecutionContext = await workflowExecutionContext.CreateActivityExecutionContextAsync(activity);
 141158            activityExecutionContext.Id = activityExecutionContextState.Id;
 141159            activityExecutionContext.CallStackDepth = activityExecutionContextState.CallStackDepth;
 141160            activityExecutionContext.Properties.Merge(properties);
 141161            activityExecutionContext.Metadata.Merge(metadata);
 162
 141163            if(activityExecutionContextState.ActivityState != null)
 141164                activityExecutionContext.ActivityState.Merge(activityExecutionContextState.ActivityState);
 165
 141166            activityExecutionContext.TransitionTo(activityExecutionContextState.Status);
 141167            activityExecutionContext.IsExecuting = activityExecutionContextState.IsExecuting;
 141168            activityExecutionContext.AggregateFaultCount = activityExecutionContextState.FaultCount;
 141169            activityExecutionContext.StartedAt = activityExecutionContextState.StartedAt;
 141170            activityExecutionContext.CompletedAt = activityExecutionContextState.CompletedAt;
 141171            activityExecutionContext.Tag = activityExecutionContextState.Tag;
 141172            activityExecutionContext.DynamicVariables = activityExecutionContextState.DynamicVariables;
 173
 141174            return activityExecutionContext;
 143175        }
 141176    }
 177
 178    private void ApplyCompletionCallbacks(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
 179    {
 482180        foreach (var completionCallbackEntry in state.CompletionCallbacks)
 181        {
 392182            var ownerActivityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x
 100183            if (ownerActivityExecutionContext == null)
 184            {
 1185                LogSkippedStateReference(
 1186                    state,
 1187                    workflowExecutionContext,
 1188                    "CompletionCallbackOwner",
 1189                    completionCallbackOwnerInstanceId: completionCallbackEntry.OwnerInstanceId,
 1190                    completionCallbackChildNodeId: completionCallbackEntry.ChildNodeId);
 1191                continue;
 192            }
 193
 99194            var childNode = workflowExecutionContext.FindNodeById(completionCallbackEntry.ChildNodeId);
 195
 99196            if (childNode == null)
 197            {
 1198                LogSkippedStateReference(
 1199                    state,
 1200                    workflowExecutionContext,
 1201                    "CompletionCallbackChild",
 1202                    completionCallbackOwnerInstanceId: completionCallbackEntry.OwnerInstanceId,
 1203                    completionCallbackChildNodeId: completionCallbackEntry.ChildNodeId);
 1204                continue;
 205            }
 206
 98207            var callbackName = completionCallbackEntry.MethodName;
 98208            var callbackDelegate = !string.IsNullOrEmpty(callbackName) ? ownerActivityExecutionContext.Activity.GetActiv
 98209            var tag = completionCallbackEntry.Tag;
 98210            workflowExecutionContext.AddCompletionCallback(ownerActivityExecutionContext, childNode, callbackDelegate, t
 211        }
 141212    }
 213
 214    private void LogSkippedStateReference(
 215        WorkflowState state,
 216        WorkflowExecutionContext workflowExecutionContext,
 217        string skipKind,
 218        string? activityExecutionContextId = null,
 219        string? scheduledActivityNodeId = null,
 220        string? completionCallbackOwnerInstanceId = null,
 221        string? completionCallbackChildNodeId = null)
 222    {
 4223        var targetIdentity = workflowExecutionContext.Workflow.Identity;
 4224        var isWorkflowDefinitionVersionMigration = !string.Equals(state.DefinitionVersionId, targetIdentity.Id, StringCo
 4225        var classification = isWorkflowDefinitionVersionMigration ? "MigrationCompatible" : "Unexpected";
 226
 4227        logger.LogWarning(
 4228            "Skipping unresolved workflow state reference {WorkflowStateSkipKind} ({WorkflowStateSkipClassification}) fo
 4229            skipKind,
 4230            classification,
 4231            state.Id,
 4232            state.DefinitionId,
 4233            state.DefinitionVersionId,
 4234            state.DefinitionVersion,
 4235            targetIdentity.DefinitionId,
 4236            targetIdentity.Id,
 4237            targetIdentity.Version,
 4238            activityExecutionContextId,
 4239            scheduledActivityNodeId,
 4240            completionCallbackOwnerInstanceId,
 4241            completionCallbackChildNodeId,
 4242            isWorkflowDefinitionVersionMigration);
 4243    }
 244
 245    private void ApplyScheduledActivities(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
 246    {
 286247        foreach (var activityWorkItemState in state.ScheduledActivities)
 248        {
 2249            var activity = workflowExecutionContext.FindActivityByNodeId(activityWorkItemState.ActivityNodeId);
 250
 2251            if (activity == null)
 252                continue;
 253
 4254            var ownerContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Id == activityWo
 4255            var existingActivityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x =
 2256            var variables = activityWorkItemState.Variables;
 2257            var input = activityWorkItemState.Input;
 2258            var tag = activityWorkItemState.Tag;
 2259            var workItem = new ActivityWorkItem(
 2260                activity,
 2261                ownerContext,
 2262                tag,
 2263                variables,
 2264                existingActivityExecutionContext,
 2265                input,
 2266                activityWorkItemState.SchedulingActivityExecutionId,
 2267                activityWorkItemState.SchedulingWorkflowInstanceId,
 2268                activityWorkItemState.SchedulingCallStackDepth);
 2269            workflowExecutionContext.Scheduler.Schedule(workItem);
 270        }
 141271    }
 272
 273    private static void ExtractCompletionCallbacks(WorkflowState state, WorkflowExecutionContext workflowExecutionContex
 274    {
 275        // Assert that all referenced owner contexts exist.
 516276        var activeContexts = workflowExecutionContext.GetActiveActivityExecutionContexts().ToList();
 1306277        foreach (var completionCallback in workflowExecutionContext.CompletionCallbacks)
 278        {
 397279            var ownerContext = activeContexts.FirstOrDefault(x => x == completionCallback.Owner);
 280
 137281            if (ownerContext == null)
 0282                throw new("Lost an owner context");
 283        }
 284
 653285        var completionCallbacks = workflowExecutionContext.CompletionCallbacks.Select(x => new CompletionCallbackState(x
 286
 516287        state.CompletionCallbacks = completionCallbacks.ToList();
 516288    }
 289
 290    private static void ExtractActiveActivityExecutionContexts(WorkflowState state, WorkflowExecutionContext workflowExe
 291    {
 292        ActivityExecutionContextState CreateActivityExecutionContextState(ActivityExecutionContext activityExecutionCont
 293        {
 665294            var parentId = activityExecutionContext.ParentActivityExecutionContext?.Id;
 295
 665296            if (parentId != null)
 297            {
 531298                var parentContext = activityExecutionContext.WorkflowExecutionContext.ActivityExecutionContexts.FirstOrD
 299
 152300                if (parentContext == null)
 0301                    throw new("We lost a context. This could indicate a bug in a parent activity that completed before (
 302            }
 303
 665304            var activityExecutionContextState = new ActivityExecutionContextState
 665305            {
 665306                Id = activityExecutionContext.Id,
 665307                CallStackDepth = activityExecutionContext.CallStackDepth,
 665308                ParentContextId = activityExecutionContext.ParentActivityExecutionContext?.Id,
 665309                ScheduledActivityNodeId = activityExecutionContext.NodeId,
 665310                OwnerActivityNodeId = activityExecutionContext.ParentActivityExecutionContext?.NodeId,
 665311                Properties = activityExecutionContext.Properties,
 665312                Metadata = activityExecutionContext.Metadata,
 665313                ActivityState = activityExecutionContext.ActivityState,
 665314                Status = activityExecutionContext.Status,
 665315                IsExecuting = activityExecutionContext.IsExecuting,
 665316                FaultCount = activityExecutionContext.AggregateFaultCount,
 665317                StartedAt = activityExecutionContext.StartedAt,
 665318                CompletedAt = activityExecutionContext.CompletedAt,
 665319                Tag = activityExecutionContext.Tag,
 665320                DynamicVariables = activityExecutionContext.DynamicVariables,
 665321            };
 665322            return activityExecutionContextState;
 323        }
 324
 325        // Only persist non-completed contexts.
 516326        state.ActivityExecutionContexts = workflowExecutionContext.GetActiveActivityExecutionContexts().Reverse().Select
 516327    }
 328
 329    private void ExtractScheduledActivities(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
 330    {
 516331        var scheduledActivities = workflowExecutionContext
 516332            .Scheduler.List()
 518333            .Select(x => new ActivityWorkItemState
 518334            {
 518335                ActivityNodeId = x.Activity.NodeId,
 518336                OwnerContextId = x.Owner?.Id,
 518337                Tag = x.Tag,
 518338                Variables = x.Variables?.ToList(),
 518339                ExistingActivityExecutionContextId = x.ExistingActivityExecutionContext?.Id,
 518340                Input = x.Input,
 518341                SchedulingActivityExecutionId = x.SchedulingActivityExecutionId,
 518342                SchedulingWorkflowInstanceId = x.SchedulingWorkflowInstanceId,
 518343                SchedulingCallStackDepth = x.SchedulingCallStackDepth,
 518344            });
 345
 516346        state.ScheduledActivities = scheduledActivities.ToList();
 516347    }
 348}