diff --git a/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj b/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj index 4e4903a35..04a927c5a 100644 --- a/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj +++ b/src/apps/Elsa.Server.Web/Elsa.Server.Web.csproj @@ -59,7 +59,6 @@ - diff --git a/src/common/Elsa.Testing.Shared.Integration/DispatchWorkflowExtensions.cs b/src/common/Elsa.Testing.Shared.Integration/DispatchWorkflowExtensions.cs new file mode 100644 index 000000000..fcd9264a4 --- /dev/null +++ b/src/common/Elsa.Testing.Shared.Integration/DispatchWorkflowExtensions.cs @@ -0,0 +1,86 @@ +using Elsa.Extensions; +using Elsa.Features.Services; +using Elsa.Mediator.Contracts; +using Elsa.Workflows; +using Elsa.Workflows.Notifications; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Requests; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; + +namespace Elsa.Testing.Shared; + +public static class DispatchWorkflowExtensions +{ + public static async Task DispatchWorkflowAndRunToCompletion( + this IWorkflow workflowDefinition, + Action? configureServices = null, + Action? configureElsa = null, + string? instanceId = null, + TimeSpan? timeout = null) + { + var semaphore = new SemaphoreSlim(0, 1); + WorkflowFinished? workflowFinishedRecord = null; + + var host = Host.CreateDefaultBuilder() + .ConfigureServices(services => + { + configureServices?.Invoke(services); + + // This notification handler will capture the WorkflowFinished record (to be returned) and release the semaphore. + services.AddNotificationHandler(sp => new(notification => + { + workflowFinishedRecord = notification; + semaphore.Release(); + })); + + services.AddElsa(elsa => configureElsa?.Invoke(elsa)); + }) + .Build(); + + try + { + // Start the host. + await host.StartAsync(CancellationToken.None); + using var scope = host.Services.CreateScope(); + var serviceProvider = scope.ServiceProvider; + await serviceProvider.PopulateRegistriesAsync(); + + // Build the workflow + var workflowBuilderFactory = serviceProvider.GetRequiredService(); + var workflow = await workflowBuilderFactory.CreateBuilder().BuildWorkflowAsync(workflowDefinition); + + // Register the workflow + var workflowRegistry = serviceProvider.GetRequiredService(); + await workflowRegistry.RegisterAsync(workflow); + + // Dispatch the workflow + var workflowDispatcher = serviceProvider.GetRequiredService(); + var dispatchWorkflowResponse = await workflowDispatcher.DispatchAsync(new DispatchWorkflowDefinitionRequest + { + DefinitionVersionId = workflow.DefinitionHandle.DefinitionVersionId, + InstanceId = instanceId ?? Guid.NewGuid().ToString(), + }); + dispatchWorkflowResponse.ThrowIfFailed(); + + // Wait for the workflow to complete, and then return the WorkflowFinished notification. + var signaled = await semaphore.WaitAsync(timeout ?? TimeSpan.FromSeconds(5)); + return signaled ? workflowFinishedRecord : null; + } + finally + { + // Stop the host. + await host.StopAsync(CancellationToken.None); + } + } + + class WorkflowFinishedAction(Action action) : INotificationHandler + { + public Task HandleAsync(WorkflowFinished notification, CancellationToken cancellationToken) + { + action(notification); + return Task.CompletedTask; + } + } +} \ No newline at end of file diff --git a/src/common/Elsa.Testing.Shared.Integration/RunActivityExtensions.cs b/src/common/Elsa.Testing.Shared.Integration/RunActivityExtensions.cs new file mode 100644 index 000000000..52bc601bb --- /dev/null +++ b/src/common/Elsa.Testing.Shared.Integration/RunActivityExtensions.cs @@ -0,0 +1,52 @@ +using Elsa.Common.Models; +using Elsa.Workflows; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Models; +using Elsa.Workflows.Options; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Filters; +using Elsa.Workflows.Runtime.Messages; +using Elsa.Workflows.State; +using JetBrains.Annotations; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Testing.Shared; + +/// +/// Provides extension methods for . +/// +[PublicAPI] +public static class RunActivityExtensions +{ + /// + /// Runs the specified activity. + /// + /// The service provider. + /// The activity to run. + /// An optional cancellation token. + /// The result of running the activity. + public static async Task RunActivityAsync(this IServiceProvider services, IActivity activity, CancellationToken cancellationToken = default) + { + await services.PopulateRegistriesAsync(); + var workflowRunner = services.GetRequiredService(); + var result = await workflowRunner.RunAsync(activity, cancellationToken: cancellationToken); + return result; + } + + /// + /// Runs the specified activity. + /// + /// The service provider. + /// The activity to run. + /// An set of options. + /// An optional cancellation token. + /// The result of running the activity. + public static async Task RunActivityAsync(this IServiceProvider services, IActivity activity, RunWorkflowOptions options, CancellationToken cancellationToken = default) + { + var workflowRunner = services.GetRequiredService(); + var result = await workflowRunner.RunAsync(activity, options, cancellationToken); + return result; + } +} \ No newline at end of file diff --git a/src/common/Elsa.Testing.Shared.Integration/RunWorkflowExtensions.cs b/src/common/Elsa.Testing.Shared.Integration/RunWorkflowExtensions.cs new file mode 100644 index 000000000..25e8b3ce8 --- /dev/null +++ b/src/common/Elsa.Testing.Shared.Integration/RunWorkflowExtensions.cs @@ -0,0 +1,81 @@ +using Elsa.Common.Models; +using Elsa.Workflows; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Models; +using Elsa.Workflows.Options; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Filters; +using Elsa.Workflows.Runtime.Messages; +using Elsa.Workflows.State; +using JetBrains.Annotations; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Testing.Shared; + +/// +/// Provides extension methods for . +/// +[PublicAPI] +public static class RunWorkflowExtensions +{ + /// + /// Runs a workflow until its end, automatically resuming any bookmark it encounters. + /// + /// The services. + /// The ID of the workflow definition. + /// An optional dictionary of input values. + /// An optional set of options to specify the version of the workflow definition to retrieve. + /// The workflow state. + public static async Task RunWorkflowUntilEndAsync(this IServiceProvider services, + string workflowDefinitionId, + IDictionary? input = null, + VersionOptions? versionOptions = null) + { + var workflowDefinitionService = services.GetRequiredService(); + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, versionOptions ?? VersionOptions.Published); + var workflowRuntime = services.GetRequiredService(); + var workflowClient = await workflowRuntime.CreateClientAsync(); + var response = await workflowClient.CreateAndRunInstanceAsync(new() + { + WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionVersionId(workflowGraph!.Workflow.Identity.Id), + Input = input + }); + + var bookmarkStore = services.GetRequiredService(); + + // Continue resuming the workflow for as long as there are bookmarks to resume and the workflow is not Finished. + while (response.Status != WorkflowStatus.Finished) + { + var bookmarks = (await bookmarkStore.FindManyAsync(new() + { + WorkflowInstanceId = response.WorkflowInstanceId + })).ToList(); + + if (!bookmarks.Any()) + break; + + foreach (var bookmark in bookmarks) + { + var runRequest = new RunWorkflowInstanceRequest + { + BookmarkId = bookmark.Id + }; + response = await workflowClient.RunInstanceAsync(runRequest); + } + } + + // Return the workflow state. + return await workflowClient.ExportStateAsync(); + } + + /// + /// Runs a workflow until its end, automatically resuming any bookmark it encounters. + /// + public static async Task RunWorkflowUntilEndAsync(this IServiceProvider services, IDictionary? input = null) where TWorkflow : IWorkflow + { + var workflowDefinitionId = typeof(TWorkflow).Name; + return await services.RunWorkflowUntilEndAsync(workflowDefinitionId, input); + } +} \ No newline at end of file diff --git a/src/common/Elsa.Testing.Shared.Integration/ServiceProviderExtensions.cs b/src/common/Elsa.Testing.Shared.Integration/ServiceProviderExtensions.cs index 037086ff8..b4dfd6eb8 100644 --- a/src/common/Elsa.Testing.Shared.Integration/ServiceProviderExtensions.cs +++ b/src/common/Elsa.Testing.Shared.Integration/ServiceProviderExtensions.cs @@ -53,95 +53,6 @@ public static class ServiceProviderExtensions return result.WorkflowDefinition; } - /// - /// Runs a workflow until its end, automatically resuming any bookmark it encounters. - /// - /// The services. - /// The ID of the workflow definition. - /// An optional dictionary of input values. - /// An optional set of options to specify the version of the workflow definition to retrieve. - /// The workflow state. - public static async Task RunWorkflowUntilEndAsync(this IServiceProvider services, - string workflowDefinitionId, - IDictionary? input = default, - VersionOptions? versionOptions = default) - { - var workflowDefinitionService = services.GetRequiredService(); - var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, versionOptions ?? VersionOptions.Published); - var workflowRuntime = services.GetRequiredService(); - var workflowClient = await workflowRuntime.CreateClientAsync(); - var response = await workflowClient.CreateAndRunInstanceAsync(new CreateAndRunWorkflowInstanceRequest - { - WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionVersionId(workflowGraph!.Workflow.Identity.Id), - Input = input - }); - - var bookmarkStore = services.GetRequiredService(); - - // Continue resuming the workflow for as long as there are bookmarks to resume and the workflow is not Finished. - while (response.Status != WorkflowStatus.Finished) - { - var bookmarks = (await bookmarkStore.FindManyAsync(new BookmarkFilter - { - WorkflowInstanceId = response.WorkflowInstanceId - })).ToList(); - - if (!bookmarks.Any()) - break; - - foreach (var bookmark in bookmarks) - { - var runRequest = new RunWorkflowInstanceRequest - { - BookmarkId = bookmark.Id - }; - response = await workflowClient.RunInstanceAsync(runRequest); - } - } - - // Return the workflow state. - return await workflowClient.ExportStateAsync(); - } - - /// - /// Runs a workflow until its end, automatically resuming any bookmark it encounters. - /// - public static async Task RunWorkflowUntilEndAsync(this IServiceProvider services, IDictionary? input = default) where TWorkflow : IWorkflow - { - var workflowDefinitionId = typeof(TWorkflow).Name; - return await services.RunWorkflowUntilEndAsync(workflowDefinitionId, input); - } - - /// - /// Runs the specified activity. - /// - /// The service provider. - /// The activity to run. - /// An optional cancellation token. - /// The result of running the activity. - public static async Task RunActivityAsync(this IServiceProvider services, IActivity activity, CancellationToken cancellationToken = default) - { - await services.PopulateRegistriesAsync(); - var workflowRunner = services.GetRequiredService(); - var result = await workflowRunner.RunAsync(activity, cancellationToken: cancellationToken); - return result; - } - - /// - /// Runs the specified activity. - /// - /// The service provider. - /// The activity to run. - /// An set of options. - /// An optional cancellation token. - /// The result of running the activity. - public static async Task RunActivityAsync(this IServiceProvider services, IActivity activity, RunWorkflowOptions options, CancellationToken cancellationToken = default) - { - var workflowRunner = services.GetRequiredService(); - var result = await workflowRunner.RunAsync(activity, options, cancellationToken); - return result; - } - /// /// Retrieves a workflow definition by its ID. /// diff --git a/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs b/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs index cfc14e299..48a1118bc 100644 --- a/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs +++ b/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs @@ -26,7 +26,7 @@ public class TestApplicationBuilder public TestApplicationBuilder(ITestOutputHelper testOutputHelper) { _testOutputHelper = testOutputHelper; - _services = new ServiceCollection(); + _services = new(); _services .AddSingleton(testOutputHelper) diff --git a/src/common/Elsa.Testing.Shared.Integration/TestWorkflow.cs b/src/common/Elsa.Testing.Shared.Integration/TestWorkflow.cs index 3a74edb44..326610eb9 100644 --- a/src/common/Elsa.Testing.Shared.Integration/TestWorkflow.cs +++ b/src/common/Elsa.Testing.Shared.Integration/TestWorkflow.cs @@ -1,4 +1,4 @@ -using Elsa.Workflows; +using Elsa.Workflows; namespace Elsa.Testing.Shared; diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs index c533c2b00..b73e3758e 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs @@ -11,4 +11,11 @@ public interface IActivityExecutionMapper /// Maps an activity execution context to an activity execution record. /// Task MapAsync(ActivityExecutionContext source); + + /// + /// Retrieves a dictionary containing the persistable output of an activity execution context. + /// + /// The activity execution context to extract persistable output from. + /// A dictionary containing the persistable output. + Task> GetPersistableOutputAsync(ActivityExecutionContext context); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs index 640a07744..db98a7ac3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs @@ -124,7 +124,7 @@ public class BackgroundActivityInvokerMiddleware( continue; var output = (Output?)outputDescriptor.ValueGetter(activity); - context.Set(output, outputEntry.Value); + context.Set(output, outputEntry.Value, outputDescriptor.Name); } } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundActivityInvoker.cs b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundActivityInvoker.cs index 886fc40f4..7c1a71048 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundActivityInvoker.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundActivityInvoker.cs @@ -1,10 +1,7 @@ using System.Text.Json; using Elsa.Workflows.Management; -using Elsa.Workflows.Memory; -using Elsa.Workflows.Models; using Elsa.Workflows.Runtime.Middleware.Activities; using Elsa.Workflows.Runtime.Options; -using Elsa.Workflows.Services; using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Runtime; @@ -18,6 +15,7 @@ public class BackgroundActivityInvoker( IWorkflowDefinitionService workflowDefinitionService, IVariablePersistenceManager variablePersistenceManager, IActivityInvoker activityInvoker, + IActivityExecutionMapper activityExecutionMapper, WorkflowHeartbeatGeneratorFactory workflowHeartbeatGeneratorFactory, IServiceProvider serviceProvider, ILogger logger) @@ -68,7 +66,7 @@ public class BackgroundActivityInvoker( var completed = activityExecutionContext.GetBackgroundCompleted(); var scheduledActivities = activityExecutionContext.GetBackgroundScheduledActivities().ToList(); var workflowInstanceId = scheduledBackgroundActivity.WorkflowInstanceId; - var outputValues = ExtractActivityOutput(activityExecutionContext); + var outputValues = await activityExecutionMapper.GetPersistableOutputAsync(activityExecutionContext); var properties = new Dictionary { [scheduledActivitiesKey] = JsonSerializer.Serialize(scheduledActivities), @@ -92,37 +90,4 @@ public class BackgroundActivityInvoker( }; await bookmarkQueue.EnqueueAsync(enqueuedBookmark, cancellationToken); } - - private IDictionary ExtractActivityOutput(ActivityExecutionContext activityExecutionContext) - { - var outputDescriptors = activityExecutionContext.ActivityDescriptor.Outputs; - var outputValues = new Dictionary(); - - foreach (var outputDescriptor in outputDescriptors) - { - var output = (Output?)outputDescriptor.ValueGetter(activityExecutionContext.Activity); - - if (output == null) - continue; - - var memoryBlockReference = output.MemoryBlockReference(); - - if (!activityExecutionContext.ExpressionExecutionContext.TryGetBlock(memoryBlockReference, out var memoryBlock)) - continue; - - var variableMetadata = memoryBlock.Metadata as VariableBlockMetadata; - var driver = variableMetadata?.StorageDriverType; - - // We only capture output written to the workflow itself. Other drivers like blob storage, etc. will be ignored since the foreground context will be loading those. - if (driver != typeof(WorkflowStorageDriver) && driver != typeof(WorkflowInstanceStorageDriver) && driver != null) - continue; - - var outputValue = activityExecutionContext.Get(memoryBlockReference); - - if (outputValue != null) - outputValues[outputDescriptor.Name] = outputValue; - } - - return outputValues; - } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs index 7278e4fa1..4457225f6 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -73,33 +73,18 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper * } */ - var cancellationToken = source.WorkflowExecutionContext.CancellationToken; - var workflow = (Workflow?)source.GetAncestors().FirstOrDefault(x => x.Activity is Workflow)?.Activity ?? source.WorkflowExecutionContext.Workflow; + var cancellationToken = source.CancellationToken; + var legacyActivityPersistenceProperties = source.Activity.CustomProperties.GetValueOrDefault>(LegacyLogPersistenceModeKey, () => new Dictionary())!; var rootActivityExecutionContext = source.WorkflowExecutionContext.ActivityExecutionContexts.First(x => x.ParentActivityExecutionContext == null); + var workflow = (Workflow?)source.GetAncestors().FirstOrDefault(x => x.Activity is Workflow)?.Activity ?? source.WorkflowExecutionContext.Workflow; 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 activityPersistenceProperties = source.Activity.CustomProperties.GetValueOrDefault>(LogPersistenceConfigKey, () => new Dictionary())!; var payload = GetPayload(source); - var outputs = GetOutputs(source); + var outputs = await GetPersistableOutputAsync(source, legacyActivityPersistenceProperties, activityPersistenceProperties, activityPersistencePropertyDefault, cancellationToken); + var inputs = await GetPersistableInputAsync(source, legacyActivityPersistenceProperties, activityPersistenceProperties, activityPersistencePropertyDefault, cancellationToken); - 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 + return new() { Id = source.Id, ActivityId = source.Activity.Id, @@ -120,6 +105,69 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper }; } + public async Task> GetPersistableOutputAsync(ActivityExecutionContext context) + { + var cancellationToken = context.WorkflowExecutionContext.CancellationToken; + var legacyActivityPersistenceProperties = context.Activity.CustomProperties.GetValueOrDefault>(LegacyLogPersistenceModeKey, () => new Dictionary()); + var rootActivityExecutionContext = context.WorkflowExecutionContext.ActivityExecutionContexts.First(x => x.ParentActivityExecutionContext == null); + var workflow = (Workflow?)context.GetAncestors().FirstOrDefault(x => x.Activity is Workflow)?.Activity ?? context.WorkflowExecutionContext.Workflow; + var workflowPersistenceProperty = await GetDefaultPersistenceModeAsync(rootActivityExecutionContext.ExpressionExecutionContext, workflow.CustomProperties, () => _options.Value.LogPersistenceMode, cancellationToken); + var activityPersistencePropertyDefault = await GetDefaultPersistenceModeAsync(context.ExpressionExecutionContext, context.Activity.CustomProperties, () => workflowPersistenceProperty, cancellationToken); + var activityPersistenceProperties = context.Activity.CustomProperties.GetValueOrDefault>(LogPersistenceConfigKey, () => new Dictionary()); + + return await GetPersistableOutputAsync( + context, + legacyActivityPersistenceProperties, + activityPersistenceProperties, + activityPersistencePropertyDefault, + cancellationToken); + } + + private async Task> GetPersistableOutputAsync( + ActivityExecutionContext context, + IDictionary legacyActivityPersistenceProperties, + IDictionary activityPersistenceProperties, + LogPersistenceMode activityPersistencePropertyDefault, + CancellationToken cancellationToken) + { + var outputs = GetOutputs(context); + return await GetPersistablePropertiesAsync(context, outputs, "outputs", legacyActivityPersistenceProperties, activityPersistenceProperties, activityPersistencePropertyDefault, cancellationToken); + } + + private async Task> GetPersistableInputAsync(ActivityExecutionContext context, + IDictionary legacyActivityPersistenceProperties, + IDictionary activityPersistenceProperties, + LogPersistenceMode activityPersistencePropertyDefault, + CancellationToken cancellationToken) + { + return await GetPersistablePropertiesAsync( + context, + context.ActivityState!, + "inputs", + legacyActivityPersistenceProperties, + activityPersistenceProperties, + activityPersistencePropertyDefault, + cancellationToken); + } + + private async Task> GetPersistablePropertiesAsync( + ActivityExecutionContext context, + IDictionary state, + string key, + IDictionary legacyActivityPersistenceProperties, + IDictionary activityPersistenceProperties, + LogPersistenceMode activityPersistencePropertyDefault, + CancellationToken cancellationToken) + { + return await FilterPropertiesUsingPersistenceMode( + context.ExpressionExecutionContext, + state, + legacyActivityPersistenceProperties!.GetValueOrDefault(key, () => new Dictionary())!, + activityPersistenceProperties!.GetValueOrDefault(key, () => new Dictionary())!, + activityPersistencePropertyDefault, + cancellationToken); + } + private async Task GetDefaultPersistenceModeAsync(ExpressionExecutionContext expressionExecutionContext, IDictionary customProperties, Func defaultFactory, CancellationToken cancellationToken) { var legacyProperties = customProperties.GetValueOrDefault>(LegacyLogPersistenceModeKey, () => new Dictionary()); @@ -139,7 +187,7 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper return await EvaluateLogPersistenceConfigAsync(defaultPersistenceConfig, expressionExecutionContext, defaultFactory, cancellationToken); } - private async Task> StorePropertyUsingPersistenceMode( + private async Task> FilterPropertiesUsingPersistenceMode( ExpressionExecutionContext expressionExecutionContext, IDictionary state, IDictionary obsoletePersistenceModeConfiguration, @@ -154,16 +202,9 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper 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); - } + var mode = logPersistenceConfig != null + ? await EvaluateLogPersistenceConfigAsync(logPersistenceConfig, expressionExecutionContext, () => defaultLogPersistenceMode, cancellationToken) + : obsoletePersistenceModeConfiguration.GetValueOrDefault(propKey, () => defaultLogPersistenceMode); if (mode == LogPersistenceMode.Include || mode == LogPersistenceMode.Inherit && defaultLogPersistenceMode is LogPersistenceMode.Include or LogPersistenceMode.Inherit) result.Add(value.Key, value.Value); diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs index b25390c8a..04151f9f9 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/JsonObjectToObjectRemainsJsonObject/Tests.cs @@ -1,5 +1,4 @@ using Elsa.Testing.Shared; -using Elsa.Workflows.IntegrationTests.Scenarios.JsonObjectToObjectRemainsJsonObject.Workflows; using Microsoft.Extensions.DependencyInjection; using Xunit.Abstractions; diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Activities/SampleActivity.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Activities/SampleActivity.cs new file mode 100644 index 000000000..f6f12e5d1 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Activities/SampleActivity.cs @@ -0,0 +1,27 @@ +using Elsa.Workflows.Attributes; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.RunAsynchronousActivityOutput.Activities; + +[Activity("Elsa.Test", Kind = ActivityKind.Task)] +internal class SampleActivity : CodeActivity +{ + public SampleActivity() + { + RunAsynchronously = true; + } + + public Input? Number1 { get; set; } = null; + public Input? Number2 { get; set; } = null; + public Output? Sum { get; set; } = null; + public Output? Product { get; set; } = null; + + protected override void Execute(ActivityExecutionContext context) + { + var number1 = context.Get(Number1); + var number2 = context.Get(Number2); + + context.Set(Sum, number1 + number2); + context.Set(Product, number1 * number2); + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs new file mode 100644 index 000000000..ffe94bfab --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs @@ -0,0 +1,155 @@ +using Elsa.Extensions; +using Elsa.Testing.Shared; +using Elsa.Workflows.Activities; +using Elsa.Workflows.IntegrationTests.Scenarios.RunAsynchronousActivityOutput.Activities; +using Elsa.Workflows.Memory; +using Elsa.Workflows.Runtime.Distributed; +using Elsa.Workflows.Runtime.Stores; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.RunAsynchronousActivityOutput; + +public class Tests +{ + [Theory(DisplayName = "Activity outputs captured in activity execution record")] + [InlineData(true)] + [InlineData(false)] + public async Task ActivityOutputCaptureTest(bool runAsynchronously) + { + // Arrange + var workflow = new TestWorkflow(workflowBuilder => + { + var variable1 = new Variable(); + workflowBuilder.Root = new Sequence + { + Variables = + { + variable1 + }, + Activities = + { + new SampleActivity + { + Id = "SampleActivity1", + RunAsynchronously = runAsynchronously, + Number1 = new(4), + Number2 = new(8), + Sum = new(variable1), + Product = null + } + } + }; + }); + + var activityExecutionStore = new MemoryActivityExecutionStore(new()); + + // Act + var workflowFinishedRecord = await workflow.DispatchWorkflowAndRunToCompletion( + configureElsa: elsa => + { + elsa.UseWorkflowRuntime(workflowRuntime => + { + workflowRuntime.ActivityExecutionLogStore = sp => activityExecutionStore; + }); + } + ); + + // Assert + Assert.NotNull(workflowFinishedRecord); + Assert.Equal(WorkflowStatus.Finished, workflowFinishedRecord.WorkflowState.Status); + Assert.Equal(WorkflowSubStatus.Finished, workflowFinishedRecord.WorkflowState.SubStatus); + + var activityExecutionRecord = await activityExecutionStore.FindAsync(new() + { + ActivityId = "SampleActivity1" + }); + Assert.NotNull(activityExecutionRecord?.Outputs); + Assert.Equal(2, activityExecutionRecord.Outputs!.Count); + Assert.Equal(12, activityExecutionRecord.Outputs!.GetValue("Sum")); + Assert.Equal(32, activityExecutionRecord.Outputs!.GetValue("Product")); + + var activityOutputRegister = workflowFinishedRecord.WorkflowExecutionContext.GetActivityOutputRegister(); + Assert.Equal(12, activityOutputRegister.FindOutputByActivityId("SampleActivity1", "Sum")); + Assert.Equal(32, activityOutputRegister.FindOutputByActivityId("SampleActivity1", "Product")); + } + + [Theory(DisplayName = "Activity outputs captured in activity execution record")] + [InlineData(true)] + [InlineData(false)] + public async Task ActivityOutputCaptureParallelTest(bool runAsynchronously) + { + // Arrange + var workflow = new TestWorkflow(workflowBuilder => + { + var variable1 = new Variable(); + var variable2 = new Variable(); + workflowBuilder.Root = new Elsa.Workflows.Activities.Parallel + { + Variables = + { + variable1, + variable2, + }, + Activities = + { + new SampleActivity + { + Id = "SampleActivity1", + RunAsynchronously = runAsynchronously, + Number1 = new(4), + Number2 = new(8), + Sum = new(variable1), + }, + new SampleActivity + { + Id = "SampleActivity2", + RunAsynchronously = runAsynchronously, + Number1 = new(2), + Number2 = new(7), + Product = new(variable2), + } + } + }; + }); + + var activityExecutionStore = new MemoryActivityExecutionStore(new()); + + // Act + var workflowFinishedRecord = await workflow.DispatchWorkflowAndRunToCompletion( + configureServices: services => + { + services.AddScoped(); + }, + configureElsa: elsa => + { + elsa.UseWorkflowRuntime(workflowRuntime => + { + workflowRuntime.ActivityExecutionLogStore = sp => activityExecutionStore; + workflowRuntime.WorkflowRuntime = sp => sp.GetRequiredService(); + }); + }); + + // Assert + Assert.NotNull(workflowFinishedRecord); + Assert.Equal(WorkflowStatus.Finished, workflowFinishedRecord.WorkflowState.Status); + Assert.Equal(WorkflowSubStatus.Finished, workflowFinishedRecord.WorkflowState.SubStatus); + + var activityExecutionRecord1 = await activityExecutionStore.FindAsync(new() + { + ActivityId = "SampleActivity1" + }); + Assert.NotNull(activityExecutionRecord1?.Outputs); + Assert.Equal(2, activityExecutionRecord1.Outputs!.Count); + Assert.Equal(12, activityExecutionRecord1.Outputs!.GetValue("Sum")); + Assert.Equal(32, activityExecutionRecord1.Outputs!.GetValue("Product")); + + var activityExecutionRecord2 = await activityExecutionStore.FindAsync(new() + { + ActivityId = "SampleActivity2" + }); + Assert.NotNull(activityExecutionRecord2?.Outputs); + Assert.Equal(2, activityExecutionRecord2.Outputs!.Count); + Assert.Equal(9, activityExecutionRecord2.Outputs!.GetValue("Sum")); + Assert.Equal(14, activityExecutionRecord2.Outputs!.GetValue("Product")); + } +} \ No newline at end of file