Adds Publish Event Activity Tests (#7093)

* Refactor workflow instance deletion to use `IWorkflowRuntime` for enhanced coordination and separation of concerns.

* Remove `EnumerableTypeConverter` and update related usages for serialization.

- Deleted the `EnumerableTypeConverter` class and its JSON serialization logic.
- Removed associated type descriptor attribute in `DefaultFormattersFeature`.
- Updated `ObjectFormatter` to handle collection serialization directly with JSON.

* Remove `EnumerableTypeConverter` tests and consolidate serialization logic into `ObjectFormatter`.

- Deleted `EnumerableTypeConverterTests` as the related functionality was removed.
- Added comprehensive tests in `ObjectFormatterTests` to handle serialization of collections and arrays with JSON.

* Add integration tests for `TriggerIndexer` to handle workflows with failing materialization

- Introduced comprehensive test scenarios verifying `DeleteTriggersAsync` behavior when workflows fail to load or partially succeed.
- Enhanced error handling in `TriggerIndexer` to skip failed workflows while ensuring remaining workflows are processed.

* Add exception handling in `TriggerIndexer.DeleteTriggersAsync` and integration tests

- Enhanced `DeleteTriggersAsync` with exception handling to skip failed workflows while processing others.
- Logged warnings for failed workflows without halting execution.
- Added comprehensive integration tests to verify behavior across success, failure, and mixed scenarios.
- Refactored tests for improved clarity, maintainability, and consistency.

* Add exception handling for `ResumeWorkflowTask` to skip deleted workflow instances

- Enhanced `ResumeWorkflowTask.ExecuteAsync` to handle `WorkflowInstanceNotFoundException` gracefully when a scheduled workflow instance is missing.
- Logged warnings for skipped executions to improve observability.
- Ensured remaining workflows and scheduled tasks are processed seamlessly without disruption.

* Add thread safety to `LocalScheduler` to prevent race conditions during concurrent scheduling

- Introduced a `lock` object to synchronize access to internal dictionaries.
- Resolved `IndexOutOfRangeException` caused by concurrent modifications during startup.
- Ensured thread-safe operations in `ScheduleAsync`, `ClearScheduleAsync`, and related methods.
- Improved reliability and stability of scheduling under concurrent workloads.

* Improve exception handling, thread safety, and workflow instance deletion

- Added exception handling in `TriggerIndexer.DeleteTriggersAsync` to skip failed workflows while continuing processing.
- Enhanced `ResumeWorkflowTask` to handle missing workflow instances gracefully and log warnings.
- Introduced thread synchronization in `LocalScheduler` with `lock` to prevent concurrent access issues.
- Implemented and refactored tests to ensure behavior consistency and improve maintainability.
- Added component tests for workflow deletion scenarios, covering running, completed, and non-existent workflows.

* Add component tests for workflow instance deletion and refactor bulk delete logic

- Added comprehensive component tests for workflow instance deletion scenarios (running, completed, bulk, and non-existent instances).
- Refactored `BulkDelete` API to use `IWorkflowInstanceManager` for proper cleanup of related records (execution logs, activity executions, bookmarks).

* Add integration tests and fakes for `TriggerIndexer` to verify behavior with failing and successful workflows

- Introduced `FailingMaterializer` and `WorkingMaterializer` for simulating failing and successful workflow materializations.
- Added `TriggerDeletionTestScenario`, `TriggerTestDataBuilder`, and related test data classes to define comprehensive test cases.
- Updated `DeleteTriggersAsync` tests with scenarios for materialization failures and mixed success.
- Improved test coverage and maintainability with reusable test data builders and utilities.

* Refactor `ActivityExecutionContextExtensions` to use instance methods for improved readability and encapsulation

* Refactor extension methods to use instance methods for improved encapsulation and readability in core workflow modules

* Add component tests for event-based workflows and update usages of `Event` activity

- Added `BlockingEventWorkflow` and `TriggerEventWorkflow` for testing event-based workflow scenarios.
- Added `EventTests` to verify workflow behavior with event publishing and triggering.
- Refactored existing integration tests to use `Runtime.Activities.Event` for consistency.

* Add unit tests for `EventBase` functionality

- Introduced `EventBaseTests` to validate core `EventBase` logic, including bookmark creation, event stimulus handling, and callback invocation.
- Added tests for scenarios involving event payloads, trigger indexing, and result output determination.
- Verified behavior consistency with various event names and callback executions.

* Add tests and workflows to validate event publishing and consumption

- Introduced `ConsumerWorkflow`, `PublishGlobalEventWorkflow`, and `PublishAndConsumeEventWorkflow` to test global and local event publishing scenarios.
- Added component tests (`PublishEventTests`) to verify event propagation and workflow triggering mechanisms.
- Implemented unit tests for `PublishEvent` with various parameters (event name, payload, correlation ID).

* Remove unused `using` directives in event-related component tests and workflows

* Refactor `PublishEventTests` and `EventBaseTests` to improve test coverage, simplify test logic, and consolidate duplicate code.

* Add `NullIfWhiteSpace` extension method and update `PublishEvent` logic to use it in correlation ID handling

- Refactored `PublishEventTests` to account for cases where correlation ID is whitespace.
- Improved test coverage for `PublishEvent` activity with additional inline test cases.

* Refactor `PublishEventTests` to verify payload transmission and enhance `ConsumerWorkflow` to capture and validate event payloads.

* Refactor `PublishEventTests` to add timeout mechanism for workflow instance retrieval; enhance `ConsumerWorkflow` to declare output variable for payload validation.

* Update test/component/Elsa.Workflows.ComponentTests/Scenarios/Activities/Primitives/Event/PublishEventTests.cs

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

* Remove `EventBaseTests` and `CancelInboundAncestorsAsync` for cleanup and redundant logic removal.

---------

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
Sipke Schoorstra 2025-11-25 19:06:16 +01:00 committed by GitHub
parent a5a9597ef3
commit faebea76a0
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 393 additions and 3 deletions

View file

@ -6,4 +6,5 @@ public static class StringExtensions
public static string WithDefault(this string? value, string defaultValue) => !string.IsNullOrWhiteSpace(value) ? value : defaultValue;
public static string EmptyIfNull(this string? value) => value ?? "";
public static string? NullIfEmpty(this string? value) => value == "" ? null : value;
public static string? NullIfWhiteSpace(this string? value) => string.IsNullOrWhiteSpace(value) ? null : value;
}

View file

@ -41,7 +41,7 @@ public class PublishEvent([CallerFilePath] string? source = null, [CallerLineNum
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
var eventName = EventName.Get(context);
var correlationId = CorrelationId.GetOrDefault(context).NullIfEmpty();
var correlationId = CorrelationId.GetOrDefault(context).NullIfWhiteSpace();
var isLocalEvent = IsLocalEvent.GetOrDefault(context);
var workflowInstanceId = isLocalEvent ? context.WorkflowExecutionContext.Id : null;
var payload = Payload.GetOrDefault(context);

View file

@ -0,0 +1,132 @@
using System.Text.Json;
using Elsa.Common.Models;
using Elsa.Testing.Shared;
using Elsa.Testing.Shared.Services;
using Elsa.Workflows.ComponentTests.Abstractions;
using Elsa.Workflows.ComponentTests.Fixtures;
using Elsa.Workflows.ComponentTests.Scenarios.Activities.Primitives.Event.Workflows;
using Elsa.Workflows.Management;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime;
using Elsa.Workflows.Runtime.Messages;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.ComponentTests.Scenarios.Activities.Primitives.Event;
public class PublishEventTests : AppComponentTest
{
private readonly AsyncWorkflowRunner _workflowRunner;
private readonly IWorkflowInstanceStore _workflowInstanceStore;
private readonly IWorkflowRuntime _workflowRuntime;
private readonly WorkflowEvents _workflowEvents;
public PublishEventTests(App app) : base(app)
{
_workflowRunner = Scope.ServiceProvider.GetRequiredService<AsyncWorkflowRunner>();
_workflowInstanceStore = Scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
_workflowRuntime = Scope.ServiceProvider.GetRequiredService<IWorkflowRuntime>();
_workflowEvents = Scope.ServiceProvider.GetRequiredService<WorkflowEvents>();
}
[Fact]
public async Task PublishEvent_LocalEvent_CompletesWorkflow()
{
// Act
var result = await _workflowRunner.RunAndAwaitWorkflowCompletionAsync(WorkflowDefinitionHandle.ByDefinitionId(PublishAndConsumeEventWorkflow.DefinitionId, VersionOptions.Published));
// Assert
Assert.Equal(WorkflowSubStatus.Finished, result.WorkflowExecutionContext.SubStatus);
}
[Fact]
public async Task PublishEvent_GlobalEvent_TriggersConsumerWorkflow()
{
// Arrange
var correlationId = Guid.NewGuid().ToString();
// Act - Publish global event
await RunWorkflowAsync(PublishGlobalEventWorkflow.DefinitionId, correlationId);
// Assert - Consumer workflow was triggered and completed
var consumerInstance = await GetSingleWorkflowInstanceAsync(ConsumerWorkflow.DefinitionId, correlationId);
Assert.Equal(WorkflowStatus.Finished, consumerInstance.Status);
Assert.Equal(WorkflowSubStatus.Finished, consumerInstance.SubStatus);
}
[Fact]
public async Task PublishEvent_WithPayload_TransmitsPayloadToConsumer()
{
// Arrange
var correlationId = Guid.NewGuid().ToString();
// Act - Publish global event with payload
await RunWorkflowAsync(PublishGlobalEventWorkflow.DefinitionId, correlationId);
// Assert - Consumer workflow received the payload
var consumerInstance = await GetSingleWorkflowInstanceAsync(ConsumerWorkflow.DefinitionId, correlationId);
Assert.Equal(WorkflowStatus.Finished, consumerInstance.Status);
Assert.Equal(WorkflowSubStatus.Finished, consumerInstance.SubStatus);
// Verify the payload was captured in the output
Assert.True(consumerInstance.WorkflowState.Output.TryGetValue("ReceivedPayload", out var receivedPayload), "Consumer workflow should have ReceivedPayload output");
Assert.NotNull(receivedPayload);
// Verify the payload structure and content
var payloadJson = JsonSerializer.Serialize(receivedPayload);
Assert.Contains("\"Status\"", payloadJson);
Assert.Contains("\"Shipped\"", payloadJson);
}
private async Task RunWorkflowAsync(string definitionId, string? correlationId = null)
{
var workflowClient = await _workflowRuntime.CreateClientAsync();
await workflowClient.CreateInstanceAsync(new()
{
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Published),
CorrelationId = correlationId
});
await workflowClient.RunInstanceAsync(RunWorkflowInstanceRequest.Empty);
}
private async Task<WorkflowInstance> GetSingleWorkflowInstanceAsync(string definitionId, string correlationId, int timeoutMs = 5000)
{
var tcs = new TaskCompletionSource<WorkflowInstance>();
var cts = new CancellationTokenSource(timeoutMs);
// Register cancellation to fail the task on timeout
cts.Token.Register(() => tcs.TrySetException(new TimeoutException($"Workflow instance with DefinitionId '{definitionId}' and CorrelationId '{correlationId}' was not saved within {timeoutMs}ms")));
// Subscribe to the WorkflowInstanceSaved event
void OnWorkflowInstanceSaved(object? sender, WorkflowInstanceSavedEventArgs args)
{
if (args.WorkflowInstance.DefinitionId == definitionId && args.WorkflowInstance.CorrelationId == correlationId)
{
tcs.TrySetResult(args.WorkflowInstance);
}
}
_workflowEvents.WorkflowInstanceSaved += OnWorkflowInstanceSaved;
try
{
// Check if the instance already exists in the database
var existingInstances = (await _workflowInstanceStore.FindManyAsync(new()
{
DefinitionId = definitionId,
CorrelationId = correlationId
}, cts.Token)).ToList();
if (existingInstances.Any())
return Assert.Single(existingInstances);
// Wait for the event to be raised
return await tcs.Task;
}
finally
{
_workflowEvents.WorkflowInstanceSaved -= OnWorkflowInstanceSaved;
cts.Dispose();
}
}
}

View file

@ -0,0 +1,42 @@
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Activities.SetOutput;
using Elsa.Workflows.Memory;
namespace Elsa.Workflows.ComponentTests.Scenarios.Activities.Primitives.Event.Workflows;
/// <summary>
/// A workflow that listens for global events as a trigger and captures the payload.
/// </summary>
public class ConsumerWorkflow : WorkflowBase
{
public static readonly string DefinitionId = Guid.NewGuid().ToString();
protected override void Build(IWorkflowBuilder workflow)
{
var eventPayload = new Variable<object>();
workflow.WithDefinitionId(DefinitionId);
workflow.WithVariable(eventPayload);
workflow.WithOutput<object>("ReceivedPayload"); // Declare the output
var eventActivity = new Runtime.Activities.Event("GlobalOrderEvent")
{
CanStartWorkflow = true,
Result = new(eventPayload)
};
workflow.Root = new Sequence
{
Activities =
{
eventActivity,
new SetOutput
{
OutputName = new("ReceivedPayload"),
OutputValue = new(eventPayload)
},
new End()
}
};
}
}

View file

@ -0,0 +1,34 @@
using Elsa.Workflows.Activities;
using Elsa.Workflows.Runtime.Activities;
namespace Elsa.Workflows.ComponentTests.Scenarios.Activities.Primitives.Event.Workflows;
/// <summary>
/// A workflow that publishes an event to itself (local event).
/// </summary>
public class PublishAndConsumeEventWorkflow : WorkflowBase
{
public static readonly string DefinitionId = Guid.NewGuid().ToString();
protected override void Build(IWorkflowBuilder workflow)
{
workflow.WithDefinitionId(DefinitionId);
workflow.Root = new Sequence
{
Activities =
{
new Start(),
// Publish a local event
new PublishEvent
{
EventName = new("LocalOrderEvent"),
IsLocalEvent = new(true),
Payload = new(new { OrderId = 123 })
},
// Wait for the local event
new Elsa.Workflows.Runtime.Activities.Event("LocalOrderEvent"),
new End()
}
};
}
}

View file

@ -0,0 +1,33 @@
using Elsa.Extensions;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Runtime.Activities;
namespace Elsa.Workflows.ComponentTests.Scenarios.Activities.Primitives.Event.Workflows;
/// <summary>
/// A workflow that publishes a global event with the workflow's correlation ID.
/// </summary>
public class PublishGlobalEventWorkflow : WorkflowBase
{
public static readonly string DefinitionId = Guid.NewGuid().ToString();
protected override void Build(IWorkflowBuilder workflow)
{
workflow.WithDefinitionId(DefinitionId);
workflow.Root = new Sequence
{
Activities =
{
new Start(),
new PublishEvent
{
EventName = new("GlobalOrderEvent"),
CorrelationId = new(context => context.GetWorkflowExecutionContext().CorrelationId),
IsLocalEvent = new(false),
Payload = new(new { Status = "Shipped" })
},
new End()
}
};
}
}

View file

@ -6,7 +6,7 @@ using Elsa.Workflows.Runtime;
using Elsa.Workflows.Runtime.Activities;
using Elsa.Workflows.Runtime.Stimuli;
namespace Elsa.Activities.UnitTests.Event;
namespace Elsa.Activities.UnitTests.Primitives;
public class EventBaseTests
{
@ -64,7 +64,7 @@ public class EventBaseTests
// Act
var context = await ExecuteAsync(activity);
context.WorkflowExecutionContext.Input[Elsa.Workflows.Runtime.Activities.Event.EventInputWorkflowInputKey] = expectedInput;
context.WorkflowExecutionContext.Input[Event.EventInputWorkflowInputKey] = expectedInput;
await activity.InvokeCallbackAsync(context);
// Assert

View file

@ -0,0 +1,148 @@
using Elsa.Testing.Shared;
using Elsa.Workflows;
using Elsa.Workflows.Runtime;
using Elsa.Workflows.Runtime.Activities;
using Microsoft.Extensions.DependencyInjection;
using NSubstitute;
namespace Elsa.Activities.UnitTests.Primitives;
public class PublishEventTests
{
[Theory]
[InlineData("OrderCreated", null, null, false)]
[InlineData("OrderCreated", "correlation-123", "correlation-123", false)]
[InlineData("OrderEvent", "", null, true)]
[InlineData("OrderEvent", " ", null, true)]
public async Task ExecuteAsync_PublishesEvent_WithParameters(string eventName, string? correlationId, string? expectedCorrelationId, bool expectNullCorrelation)
{
// Arrange
var publisher = Substitute.For<IEventPublisher>();
// Act
await ExecuteAsync(CreateActivity(eventName, correlationId), publisher);
// Assert
await AssertPublishedAsync(
publisher,
eventName,
correlationId: expectedCorrelationId,
expectNullCorrelationId: expectNullCorrelation);
}
[Theory]
[InlineData(true, false)]
[InlineData(false, true)]
public async Task ExecuteAsync_LocalEvent_PassesCorrectWorkflowInstanceId(bool isLocalEvent, bool expectNull)
{
// Arrange
const string eventName = "TestEvent";
var publisher = Substitute.For<IEventPublisher>();
// Act
var context = await ExecuteAsync(CreateActivity(eventName, isLocalEvent: isLocalEvent), publisher);
// Assert
await AssertPublishedAsync(
publisher,
eventName,
workflowInstanceId: expectNull ? null : context.WorkflowExecutionContext.Id,
expectNullWorkflowInstanceId: expectNull);
}
[Fact]
public async Task ExecuteAsync_PublishesEvent_WithPayload()
{
// Arrange
const string eventName = "OrderCreated";
var payload = new { OrderId = 123, Amount = 99.99m };
var publisher = Substitute.For<IEventPublisher>();
// Act
await ExecuteAsync(CreateActivity(eventName, payload: payload), publisher);
// Assert
await AssertPublishedAsync(publisher, eventName, payload: payload);
}
[Fact]
public async Task ExecuteAsync_CompletesActivity()
{
// Arrange
var publisher = Substitute.For<IEventPublisher>();
// Act
var context = await ExecuteAsync(CreateActivity("TestEvent"), publisher);
// Assert
Assert.Equal(ActivityStatus.Completed, context.Status);
}
[Fact]
public async Task ExecuteAsync_WithAllParameters_PassesAllValuesToPublisher()
{
// Arrange
const string eventName = "CompleteOrderEvent";
const string correlationId = "correlation-456";
var payload = new { Status = "Shipped" };
var publisher = Substitute.For<IEventPublisher>();
// Act
var context = await ExecuteAsync(CreateActivity(eventName, correlationId, isLocalEvent: true, payload: payload), publisher);
// Assert
await publisher.Received(1).PublishAsync(
eventName,
correlationId,
context.WorkflowExecutionContext.Id,
null,
payload,
true,
Arg.Any<CancellationToken>());
}
private static PublishEvent CreateActivity(
string eventName,
string? correlationId = null,
bool? isLocalEvent = null,
object? payload = null) =>
new()
{
EventName = new(eventName),
CorrelationId = correlationId != null ? new(correlationId) : null!,
IsLocalEvent = isLocalEvent.HasValue ? new(isLocalEvent.Value) : null!,
Payload = payload != null ? new(payload) : null!
};
private static async Task<ActivityExecutionContext> ExecuteAsync(PublishEvent activity, IEventPublisher publisher) =>
await new ActivityTestFixture(activity)
.ConfigureServices(services => services.AddSingleton(publisher))
.ExecuteAsync();
private static async Task AssertPublishedAsync(
IEventPublisher publisher,
string eventName,
string? correlationId = null,
string? workflowInstanceId = null,
object? payload = null,
bool expectNullCorrelationId = false,
bool expectNullWorkflowInstanceId = false)
{
var correlationIdArg = expectNullCorrelationId
? Arg.Is<string?>(x => x == null)
: correlationId ?? Arg.Any<string?>();
var workflowInstanceIdArg = expectNullWorkflowInstanceId
? Arg.Is<string?>(x => x == null)
: workflowInstanceId ?? Arg.Any<string?>();
await publisher.Received(1).PublishAsync(
eventName,
correlationIdArg,
workflowInstanceIdArg,
Arg.Any<string?>(),
payload ?? Arg.Any<object?>(),
Arg.Any<bool>(),
Arg.Any<CancellationToken>());
}
}