Fix Output Persistence of Async Activities (#6542)
* Updated output handling of asynchronously run workflows to be the same as when run synchronously * updated tests * Refactor activity execution mapping and output persistence Introduced `GetPersistableOutputAsync` in `IActivityExecutionMapper` to streamline output persistence logic. Refactored the handling of activity persistence properties, replacing repetitive code with reusable methods. Removed unused dependencies and redundant methods, optimizing code readability and maintainability. * Remove docker-compose-datadog.yml from solution file. The docker-compose-datadog.yml file is no longer included in the solution structure. This change cleans up unused references to ensure the solution remains consistent and up-to-date. * Refactor workflow extensions and add new utilities Split and reorganize workflow-related extension methods into `RunActivityExtensions` and `RunWorkflowExtensions` for better modularity. Removed deprecated methods from `ServiceProviderExtensions`. Updated tests and usages to reflect these changes. --------- Co-authored-by: Bob Hauser <rhauser@kinaxis.com>
This commit is contained in:
parent
f7743a0fe6
commit
7fadf51be6
|
|
@ -59,7 +59,6 @@
|
|||
<ItemGroup>
|
||||
<PackageReference Include="Azure.Identity" />
|
||||
<PackageReference Include="Bogus" />
|
||||
<PackageReference Include="Datadog.Trace.Bundle" />
|
||||
<PackageReference Include="DistributedLock.Postgres" />
|
||||
<PackageReference Include="DistributedLock.Redis" />
|
||||
<PackageReference Include="FluentStorage.Azure.Blobs" />
|
||||
|
|
|
|||
|
|
@ -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<WorkflowFinished?> DispatchWorkflowAndRunToCompletion(
|
||||
this IWorkflow workflowDefinition,
|
||||
Action<IServiceCollection>? configureServices = null,
|
||||
Action<IModule>? 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<WorkflowFinishedAction, WorkflowFinished>(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<IWorkflowBuilderFactory>();
|
||||
var workflow = await workflowBuilderFactory.CreateBuilder().BuildWorkflowAsync(workflowDefinition);
|
||||
|
||||
// Register the workflow
|
||||
var workflowRegistry = serviceProvider.GetRequiredService<IWorkflowRegistry>();
|
||||
await workflowRegistry.RegisterAsync(workflow);
|
||||
|
||||
// Dispatch the workflow
|
||||
var workflowDispatcher = serviceProvider.GetRequiredService<IWorkflowDispatcher>();
|
||||
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<WorkflowFinished> action) : INotificationHandler<WorkflowFinished>
|
||||
{
|
||||
public Task HandleAsync(WorkflowFinished notification, CancellationToken cancellationToken)
|
||||
{
|
||||
action(notification);
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Provides extension methods for <see cref="IServiceProvider"/>.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
public static class RunActivityExtensions
|
||||
{
|
||||
/// <summary>
|
||||
/// Runs the specified activity.
|
||||
/// </summary>
|
||||
/// <param name="services">The service provider.</param>
|
||||
/// <param name="activity">The activity to run.</param>
|
||||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
/// <returns>The result of running the activity.</returns>
|
||||
public static async Task<RunWorkflowResult> RunActivityAsync(this IServiceProvider services, IActivity activity, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await services.PopulateRegistriesAsync();
|
||||
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
|
||||
var result = await workflowRunner.RunAsync(activity, cancellationToken: cancellationToken);
|
||||
return result;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Runs the specified activity.
|
||||
/// </summary>
|
||||
/// <param name="services">The service provider.</param>
|
||||
/// <param name="activity">The activity to run.</param>
|
||||
/// <param name="options">An set of options.</param>
|
||||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
/// <returns>The result of running the activity.</returns>
|
||||
public static async Task<RunWorkflowResult> RunActivityAsync(this IServiceProvider services, IActivity activity, RunWorkflowOptions options, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
|
||||
var result = await workflowRunner.RunAsync(activity, options, cancellationToken);
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Provides extension methods for <see cref="IServiceProvider"/>.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
public static class RunWorkflowExtensions
|
||||
{
|
||||
/// <summary>
|
||||
/// Runs a workflow until its end, automatically resuming any bookmark it encounters.
|
||||
/// </summary>
|
||||
/// <param name="services">The services.</param>
|
||||
/// <param name="workflowDefinitionId">The ID of the workflow definition.</param>
|
||||
/// <param name="input">An optional dictionary of input values.</param>
|
||||
/// <param name="versionOptions">An optional set of options to specify the version of the workflow definition to retrieve.</param>
|
||||
/// <returns>The workflow state.</returns>
|
||||
public static async Task<WorkflowState> RunWorkflowUntilEndAsync(this IServiceProvider services,
|
||||
string workflowDefinitionId,
|
||||
IDictionary<string, object>? input = null,
|
||||
VersionOptions? versionOptions = null)
|
||||
{
|
||||
var workflowDefinitionService = services.GetRequiredService<IWorkflowDefinitionService>();
|
||||
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, versionOptions ?? VersionOptions.Published);
|
||||
var workflowRuntime = services.GetRequiredService<IWorkflowRuntime>();
|
||||
var workflowClient = await workflowRuntime.CreateClientAsync();
|
||||
var response = await workflowClient.CreateAndRunInstanceAsync(new()
|
||||
{
|
||||
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionVersionId(workflowGraph!.Workflow.Identity.Id),
|
||||
Input = input
|
||||
});
|
||||
|
||||
var bookmarkStore = services.GetRequiredService<IBookmarkStore>();
|
||||
|
||||
// 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();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Runs a workflow until its end, automatically resuming any bookmark it encounters.
|
||||
/// </summary>
|
||||
public static async Task<WorkflowState> RunWorkflowUntilEndAsync<TWorkflow>(this IServiceProvider services, IDictionary<string, object>? input = null) where TWorkflow : IWorkflow
|
||||
{
|
||||
var workflowDefinitionId = typeof(TWorkflow).Name;
|
||||
return await services.RunWorkflowUntilEndAsync(workflowDefinitionId, input);
|
||||
}
|
||||
}
|
||||
|
|
@ -53,95 +53,6 @@ public static class ServiceProviderExtensions
|
|||
return result.WorkflowDefinition;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Runs a workflow until its end, automatically resuming any bookmark it encounters.
|
||||
/// </summary>
|
||||
/// <param name="services">The services.</param>
|
||||
/// <param name="workflowDefinitionId">The ID of the workflow definition.</param>
|
||||
/// <param name="input">An optional dictionary of input values.</param>
|
||||
/// <param name="versionOptions">An optional set of options to specify the version of the workflow definition to retrieve.</param>
|
||||
/// <returns>The workflow state.</returns>
|
||||
public static async Task<WorkflowState> RunWorkflowUntilEndAsync(this IServiceProvider services,
|
||||
string workflowDefinitionId,
|
||||
IDictionary<string, object>? input = default,
|
||||
VersionOptions? versionOptions = default)
|
||||
{
|
||||
var workflowDefinitionService = services.GetRequiredService<IWorkflowDefinitionService>();
|
||||
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, versionOptions ?? VersionOptions.Published);
|
||||
var workflowRuntime = services.GetRequiredService<IWorkflowRuntime>();
|
||||
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<IBookmarkStore>();
|
||||
|
||||
// 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();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Runs a workflow until its end, automatically resuming any bookmark it encounters.
|
||||
/// </summary>
|
||||
public static async Task<WorkflowState> RunWorkflowUntilEndAsync<TWorkflow>(this IServiceProvider services, IDictionary<string, object>? input = default) where TWorkflow : IWorkflow
|
||||
{
|
||||
var workflowDefinitionId = typeof(TWorkflow).Name;
|
||||
return await services.RunWorkflowUntilEndAsync(workflowDefinitionId, input);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Runs the specified activity.
|
||||
/// </summary>
|
||||
/// <param name="services">The service provider.</param>
|
||||
/// <param name="activity">The activity to run.</param>
|
||||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
/// <returns>The result of running the activity.</returns>
|
||||
public static async Task<RunWorkflowResult> RunActivityAsync(this IServiceProvider services, IActivity activity, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await services.PopulateRegistriesAsync();
|
||||
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
|
||||
var result = await workflowRunner.RunAsync(activity, cancellationToken: cancellationToken);
|
||||
return result;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Runs the specified activity.
|
||||
/// </summary>
|
||||
/// <param name="services">The service provider.</param>
|
||||
/// <param name="activity">The activity to run.</param>
|
||||
/// <param name="options">An set of options.</param>
|
||||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
/// <returns>The result of running the activity.</returns>
|
||||
public static async Task<RunWorkflowResult> RunActivityAsync(this IServiceProvider services, IActivity activity, RunWorkflowOptions options, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
|
||||
var result = await workflowRunner.RunAsync(activity, options, cancellationToken);
|
||||
return result;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Retrieves a workflow definition by its ID.
|
||||
/// </summary>
|
||||
|
|
|
|||
|
|
@ -26,7 +26,7 @@ public class TestApplicationBuilder
|
|||
public TestApplicationBuilder(ITestOutputHelper testOutputHelper)
|
||||
{
|
||||
_testOutputHelper = testOutputHelper;
|
||||
_services = new ServiceCollection();
|
||||
_services = new();
|
||||
|
||||
_services
|
||||
.AddSingleton(testOutputHelper)
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
using Elsa.Workflows;
|
||||
using Elsa.Workflows;
|
||||
|
||||
namespace Elsa.Testing.Shared;
|
||||
|
||||
|
|
|
|||
|
|
@ -11,4 +11,11 @@ public interface IActivityExecutionMapper
|
|||
/// Maps an activity execution context to an activity execution record.
|
||||
/// </summary>
|
||||
Task<ActivityExecutionRecord> MapAsync(ActivityExecutionContext source);
|
||||
|
||||
/// <summary>
|
||||
/// Retrieves a dictionary containing the persistable output of an activity execution context.
|
||||
/// </summary>
|
||||
/// <param name="context">The activity execution context to extract persistable output from.</param>
|
||||
/// <returns>A dictionary containing the persistable output.</returns>
|
||||
Task<Dictionary<string, object?>> GetPersistableOutputAsync(ActivityExecutionContext context);
|
||||
}
|
||||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<BackgroundActivityInvoker> 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<string, object>
|
||||
{
|
||||
[scheduledActivitiesKey] = JsonSerializer.Serialize(scheduledActivities),
|
||||
|
|
@ -92,37 +90,4 @@ public class BackgroundActivityInvoker(
|
|||
};
|
||||
await bookmarkQueue.EnqueueAsync(enqueuedBookmark, cancellationToken);
|
||||
}
|
||||
|
||||
private IDictionary<string, object> ExtractActivityOutput(ActivityExecutionContext activityExecutionContext)
|
||||
{
|
||||
var outputDescriptors = activityExecutionContext.ActivityDescriptor.Outputs;
|
||||
var outputValues = new Dictionary<string, object>();
|
||||
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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<IDictionary<string, object?>>(LegacyLogPersistenceModeKey, () => new Dictionary<string, object?>())!;
|
||||
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<IDictionary<string, object?>>(LegacyLogPersistenceModeKey, () => new Dictionary<string, object?>());
|
||||
var activityPersistenceProperties = source.Activity.CustomProperties.GetValueOrDefault<IDictionary<string, object?>>(LogPersistenceConfigKey, () => new Dictionary<string, object?>());
|
||||
var activityPersistenceProperties = source.Activity.CustomProperties.GetValueOrDefault<IDictionary<string, object?>>(LogPersistenceConfigKey, () => new Dictionary<string, object?>())!;
|
||||
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<string, object>())!,
|
||||
activityPersistenceProperties!.GetValueOrDefault("outputs", () => new Dictionary<string, object>())!,
|
||||
activityPersistencePropertyDefault,
|
||||
cancellationToken);
|
||||
|
||||
var inputs = await StorePropertyUsingPersistenceMode(
|
||||
source.ExpressionExecutionContext,
|
||||
source.ActivityState,
|
||||
legacyActivityPersistenceProperties!.GetValueOrDefault("inputs", () => new Dictionary<string, object>())!,
|
||||
activityPersistenceProperties!.GetValueOrDefault("inputs", () => new Dictionary<string, object>())!,
|
||||
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<Dictionary<string, object?>> GetPersistableOutputAsync(ActivityExecutionContext context)
|
||||
{
|
||||
var cancellationToken = context.WorkflowExecutionContext.CancellationToken;
|
||||
var legacyActivityPersistenceProperties = context.Activity.CustomProperties.GetValueOrDefault<IDictionary<string, object?>>(LegacyLogPersistenceModeKey, () => new Dictionary<string, object?>());
|
||||
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<IDictionary<string, object?>>(LogPersistenceConfigKey, () => new Dictionary<string, object?>());
|
||||
|
||||
return await GetPersistableOutputAsync(
|
||||
context,
|
||||
legacyActivityPersistenceProperties,
|
||||
activityPersistenceProperties,
|
||||
activityPersistencePropertyDefault,
|
||||
cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<Dictionary<string, object?>> GetPersistableOutputAsync(
|
||||
ActivityExecutionContext context,
|
||||
IDictionary<string, object?> legacyActivityPersistenceProperties,
|
||||
IDictionary<string, object?> activityPersistenceProperties,
|
||||
LogPersistenceMode activityPersistencePropertyDefault,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
var outputs = GetOutputs(context);
|
||||
return await GetPersistablePropertiesAsync(context, outputs, "outputs", legacyActivityPersistenceProperties, activityPersistenceProperties, activityPersistencePropertyDefault, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<Dictionary<string, object?>> GetPersistableInputAsync(ActivityExecutionContext context,
|
||||
IDictionary<string, object?> legacyActivityPersistenceProperties,
|
||||
IDictionary<string, object?> activityPersistenceProperties,
|
||||
LogPersistenceMode activityPersistencePropertyDefault,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
return await GetPersistablePropertiesAsync(
|
||||
context,
|
||||
context.ActivityState!,
|
||||
"inputs",
|
||||
legacyActivityPersistenceProperties,
|
||||
activityPersistenceProperties,
|
||||
activityPersistencePropertyDefault,
|
||||
cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<Dictionary<string, object?>> GetPersistablePropertiesAsync(
|
||||
ActivityExecutionContext context,
|
||||
IDictionary<string, object?> state,
|
||||
string key,
|
||||
IDictionary<string, object?> legacyActivityPersistenceProperties,
|
||||
IDictionary<string, object?> activityPersistenceProperties,
|
||||
LogPersistenceMode activityPersistencePropertyDefault,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
return await FilterPropertiesUsingPersistenceMode(
|
||||
context.ExpressionExecutionContext,
|
||||
state,
|
||||
legacyActivityPersistenceProperties!.GetValueOrDefault(key, () => new Dictionary<string, object>())!,
|
||||
activityPersistenceProperties!.GetValueOrDefault(key, () => new Dictionary<string, object>())!,
|
||||
activityPersistencePropertyDefault,
|
||||
cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<LogPersistenceMode> GetDefaultPersistenceModeAsync(ExpressionExecutionContext expressionExecutionContext, IDictionary<string, object> customProperties, Func<LogPersistenceMode> defaultFactory, CancellationToken cancellationToken)
|
||||
{
|
||||
var legacyProperties = customProperties.GetValueOrDefault<IDictionary<string, object?>>(LegacyLogPersistenceModeKey, () => new Dictionary<string, object?>());
|
||||
|
|
@ -139,7 +187,7 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper
|
|||
return await EvaluateLogPersistenceConfigAsync(defaultPersistenceConfig, expressionExecutionContext, defaultFactory, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<Dictionary<string, object?>> StorePropertyUsingPersistenceMode(
|
||||
private async Task<Dictionary<string, object?>> FilterPropertiesUsingPersistenceMode(
|
||||
ExpressionExecutionContext expressionExecutionContext,
|
||||
IDictionary<string, object?> state,
|
||||
IDictionary<string, object> 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);
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
using Elsa.Testing.Shared;
|
||||
using Elsa.Workflows.IntegrationTests.Scenarios.JsonObjectToObjectRemainsJsonObject.Workflows;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Xunit.Abstractions;
|
||||
|
||||
|
|
|
|||
|
|
@ -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<int>? Number1 { get; set; } = null;
|
||||
public Input<int>? Number2 { get; set; } = null;
|
||||
public Output<int>? Sum { get; set; } = null;
|
||||
public Output<int>? 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<int>();
|
||||
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<int>("Sum"));
|
||||
Assert.Equal(32, activityExecutionRecord.Outputs!.GetValue<int>("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<int>();
|
||||
var variable2 = new Variable<int>();
|
||||
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<DistributedWorkflowRuntime>();
|
||||
},
|
||||
configureElsa: elsa =>
|
||||
{
|
||||
elsa.UseWorkflowRuntime(workflowRuntime =>
|
||||
{
|
||||
workflowRuntime.ActivityExecutionLogStore = sp => activityExecutionStore;
|
||||
workflowRuntime.WorkflowRuntime = sp => sp.GetRequiredService<DistributedWorkflowRuntime>();
|
||||
});
|
||||
});
|
||||
|
||||
// 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<int>("Sum"));
|
||||
Assert.Equal(32, activityExecutionRecord1.Outputs!.GetValue<int>("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<int>("Sum"));
|
||||
Assert.Equal(14, activityExecutionRecord2.Outputs!.GetValue<int>("Product"));
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue