diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/Composition/BulkDispatchWorkflows/BulkDispatchWorkflowsTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/Composition/BulkDispatchWorkflows/BulkDispatchWorkflowsTests.cs index 928bc317a..5f4f636b3 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/Composition/BulkDispatchWorkflows/BulkDispatchWorkflowsTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/Composition/BulkDispatchWorkflows/BulkDispatchWorkflowsTests.cs @@ -1,3 +1,4 @@ +using System.Diagnostics; using Elsa.Common.Models; using Elsa.Testing.Shared; using Elsa.Testing.Shared.Models; @@ -7,6 +8,8 @@ using Elsa.Workflows.ComponentTests.Abstractions; using Elsa.Workflows.ComponentTests.Fixtures; using Elsa.Workflows.ComponentTests.Scenarios.Activities.Composition.BulkDispatchWorkflows.Workflows; using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Models; using Elsa.Workflows.State; using Microsoft.Extensions.DependencyInjection; @@ -16,11 +19,17 @@ namespace Elsa.Workflows.ComponentTests.Scenarios.Activities.Composition.BulkDis public class BulkDispatchWorkflowsTests : AppComponentTest { private const int ChildWorkflowTimeoutSeconds = 30; + private const int InitialChildWorkflowPollingIntervalMilliseconds = 250; + private const int MaxChildWorkflowPollingIntervalMilliseconds = 1000; + // The runtime persists this marker under the activity property name; resume handlers read the same key from WorkflowState.Properties. + private const string WaitForCompletionPropertyName = nameof(Elsa.Workflows.Runtime.Activities.BulkDispatchWorkflows.WaitForCompletion); private readonly AsyncWorkflowRunner _workflowRunner; + private readonly IWorkflowInstanceStore _workflowInstanceStore; public BulkDispatchWorkflowsTests(App app) : base(app) { _workflowRunner = Scope.ServiceProvider.GetRequiredService(); + _workflowInstanceStore = Scope.ServiceProvider.GetRequiredService(); } [Fact(DisplayName = "BulkDispatchWorkflows should wait for all child workflows to complete")] @@ -37,21 +46,19 @@ public class BulkDispatchWorkflowsTests : AppComponentTest { var expectedChildCount = 3; - // Run the main workflow and wait for child workflows to complete - var (result, completedChildWorkflows) = await RunWorkflowAndWaitForChildWorkflowsAsync( - BulkDispatchFireAndForgetWorkflow.DefinitionId, + var result = await RunWorkflowAsync(BulkDispatchFireAndForgetWorkflow.DefinitionId); + + AssertWorkflowFinished(result); + var childWorkflowInstances = await WaitForChildWorkflowInstancesAsync( + result.WorkflowExecutionContext.Id, SlowBulkChildWorkflow.DefinitionId, expectedChildCount); - AssertWorkflowFinished(result); - var mainWorkflowCompletedAt = result.WorkflowExecutionContext.UpdatedAt; - - // Assert that all child workflows completed after the main workflow - Assert.Equal(expectedChildCount, completedChildWorkflows.Count); - foreach (var childContext in completedChildWorkflows) + Assert.Equal(expectedChildCount, childWorkflowInstances.Count); + foreach (var childWorkflowInstance in childWorkflowInstances) { - Assert.True(childContext.UpdatedAt > mainWorkflowCompletedAt, - $"Child workflow should complete after main workflow. Main: {mainWorkflowCompletedAt}, Child: {childContext.UpdatedAt}"); + Assert.Equal(result.WorkflowExecutionContext.Id, childWorkflowInstance.ParentWorkflowInstanceId); + Assert.False(childWorkflowInstance.WorkflowState.Properties.ContainsKey(WaitForCompletionPropertyName)); } } @@ -176,4 +183,36 @@ public class BulkDispatchWorkflowsTests : AppComponentTest workflowEvents.WorkflowStateCommitted -= OnWorkflowStateCommitted; } } + + private async Task> WaitForChildWorkflowInstancesAsync( + string parentWorkflowInstanceId, + string childWorkflowDefinitionId, + int expectedChildCount, + CancellationToken cancellationToken = default) + { + var timeout = TimeSpan.FromSeconds(ChildWorkflowTimeoutSeconds); + var pollingInterval = TimeSpan.FromMilliseconds(InitialChildWorkflowPollingIntervalMilliseconds); + var maxPollingInterval = TimeSpan.FromMilliseconds(MaxChildWorkflowPollingIntervalMilliseconds); + var stopwatch = Stopwatch.StartNew(); + var filter = new WorkflowInstanceFilter + { + DefinitionId = childWorkflowDefinitionId, + ParentWorkflowInstanceIds = [parentWorkflowInstanceId] + }; + var actualChildCount = 0; + + while (stopwatch.Elapsed < timeout) + { + var instances = (await _workflowInstanceStore.FindManyAsync(filter, cancellationToken)).ToList(); + actualChildCount = instances.Count; + + if (instances.Count >= expectedChildCount) + return instances; + + await Task.Delay(pollingInterval, cancellationToken); + pollingInterval = TimeSpan.FromMilliseconds(Math.Min(pollingInterval.TotalMilliseconds * 2, maxPollingInterval.TotalMilliseconds)); + } + + throw new TimeoutException($"Expected {expectedChildCount} child workflow instances of definition '{childWorkflowDefinitionId}' for parent '{parentWorkflowInstanceId}', but found {actualChildCount}."); + } }