From ca88051573b19f72dd4d761d990ce4eca2350784 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 19 Jan 2026 08:59:12 +0100 Subject: [PATCH] Improves workflow materializer handling (#7195) * Add tenant headers support to BackgroundWorkflowCancellationDispatcher (#7040) * Add tenant headers support to BackgroundWorkflowCancellationDispatcher * Fix 'CreateHeaders' call * Fix memory leak: Dispose IronCompressResult in Zstd codec (#7193) * Initial plan * Fix memory leak: Dispose IronCompressResult in Zstd codec and add tests Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Refactor tests to be more DRY using Theory and InlineData Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Introduce `IMaterializerRegistry` to manage workflow materializers and ensure availability checks. * Extend `IWorkflowDefinitionService` and `CachingWorkflowDefinitionService` with workflow graph lookup methods (`TryFindWorkflowGraphAsync`). Refactor caching and materialization logic for consistency. * Refactor caching interface and implementation: add `FindOrCreateAsync`, update `GetOrCreateAsync` to ensure non-null results, and improve exception handling. * Refactor `GetWorkflowGraphAsync` to use `TryFindWorkflowGraphAsync` and improve exception handling for missing workflow definitions and materializers. * Refactor caching logic to replace `GetOrCreateAsync` with `FindOrCreateAsync` for improved clarity and consistency. * Update workflow model, add event, and mark exception obsolete Updated `TimestampFilter.Column` to use a `null!` default value for clarity. Added `Event1` in the `hello-world.elsa` workflow and removed an unused folder entry from the project. Marked `WorkflowGraphNotFoundException` as obsolete with guidance to use `WorkflowDefinitionNotFoundException` instead. * Add new workflow files and exception classes for Elsa Introduced a workflow definition file "eventing.json" and new exception classes (`WorkflowDefinitionNotFoundException` and `WorkflowMaterializerNotFoundException`) to enhance handling of workflow-related errors. Also added a `WorkflowGraphFindResult` model for better workflow graph management. These changes improve the structure and functionality of the workflow system. * Add unit tests for `CachingWorkflowDefinitionService` and related helpers Introduce comprehensive unit tests to validate caching logic, workflow graph/materialization behavior, and cache key generation in `CachingWorkflowDefinitionService`. Add `WorkflowDefinitionServiceTests` and helper methods for streamlined test setup. * Enable `UseElsaScriptBlobStorage` in workflow server configuration * Refactor `BackgroundWorkflowCancellationDispatcher` to simplify object initialization and clean up XML documentation comments * Address PR #7195 review feedback: optimize caching, improve exceptions, add test coverage (#7196) * Initial plan * Apply PR review feedback: Fix exceptions, optimize caching, improve error handling Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Add unit tests for MaterializerRegistry and LocalWorkflowClient exception handling Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Add unit tests for BackgroundWorkflowCancellationDispatcher tenant headers Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Refactor `WorkflowMaterializerNotFoundException` to improve structure and usability, update related references, and simplify object initialization in test cases. * Update `WorkflowDefinitionServiceTests` to use `WorkflowMaterializerNotFoundException` in place of `InvalidOperationException` for materializer not found scenario --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> Co-authored-by: Sipke Schoorstra * Potential fix for pull request finding 'Inefficient use of ContainsKey' Co-authored-by: Copilot Autofix powered by AI <223894421+github-code-quality[bot]@users.noreply.github.com> * Refactor tests and services: simplify object initialization, use target-typed `new()` syntax, and replace `CancellationToken` with `CancellationToken.None` where applicable. * Refactor tests in `BackgroundWorkflowCancellationDispatcherTests`: improve tenant initialization and optimize header checks by replacing `TryGetValue` with `ContainsKey`. --------- Co-authored-by: Sverre Winkelmans <69142682+Sverre-W@users.noreply.github.com> Co-authored-by: Copilot <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> Co-authored-by: Copilot Autofix powered by AI <223894421+github-code-quality[bot]@users.noreply.github.com> --- .../Models/WorkflowDefinitionSummary.cs | 7 +- .../Elsa.Caching/Contracts/ICacheManager.cs | 7 +- .../Elsa.Caching/Services/CacheManager.cs | 15 +- src/modules/Elsa.Common/Codecs/Zstd.cs | 4 +- .../CachingHttpWorkflowLookupService.cs | 2 +- .../Features/ElsaScriptBlobStorageFeature.cs | 1 - .../GetByDefinitionId/Endpoint.cs | 12 +- .../Models/LinkedWorkflowDefinitionSummary.cs | 5 + .../StaticWorkflowDefinitionLinker.cs | 35 +- .../Contracts/IMaterializerRegistry.cs | 26 + .../Contracts/IWorkflowDefinitionService.cs | 59 ++ .../WorkflowMaterializerNotFoundException.cs | 6 + .../Features/WorkflowManagementFeature.cs | 1 + .../Models/TimestampFilter.cs | 2 +- .../Models/WorkflowGraphFindResult.cs | 10 + .../CachingWorkflowDefinitionService.cs | 122 +++- .../Services/MaterializerRegistry.cs | 22 + .../Services/WorkflowDefinitionService.cs | 112 +++- .../Stores/CachingWorkflowDefinitionStore.cs | 2 +- .../WorkflowDefinitionNotFoundException.cs | 8 + .../WorkflowGraphNotFoundException.cs | 1 + ...ackgroundWorkflowCancellationDispatcher.cs | 47 +- .../Services/LocalWorkflowClient.cs | 8 +- .../Stores/CachingTriggerStore.cs | 2 +- .../Elsa.Common.UnitTests/Codecs/ZstdTests.cs | 67 +++ .../Helpers/TestHelpers.cs | 40 ++ .../CachingWorkflowDefinitionServiceTests.cs | 487 ++++++++++++++++ .../Services/MaterializerRegistryTests.cs | 179 ++++++ .../WorkflowDefinitionServiceTests.cs | 524 ++++++++++++++++++ ...oundWorkflowCancellationDispatcherTests.cs | 120 ++++ ...ltWorkflowDefinitionStorePopulatorTests.cs | 40 +- .../Services/LocalWorkflowClientTests.cs | 131 +++++ 32 files changed, 1986 insertions(+), 118 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IMaterializerRegistry.cs create mode 100644 src/modules/Elsa.Workflows.Management/Exceptions/WorkflowMaterializerNotFoundException.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/WorkflowGraphFindResult.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/MaterializerRegistry.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowDefinitionNotFoundException.cs create mode 100644 test/unit/Elsa.Common.UnitTests/Codecs/ZstdTests.cs create mode 100644 test/unit/Elsa.Workflows.Management.UnitTests/Helpers/TestHelpers.cs create mode 100644 test/unit/Elsa.Workflows.Management.UnitTests/Services/CachingWorkflowDefinitionServiceTests.cs create mode 100644 test/unit/Elsa.Workflows.Management.UnitTests/Services/MaterializerRegistryTests.cs create mode 100644 test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionServiceTests.cs create mode 100644 test/unit/Elsa.Workflows.Runtime.UnitTests/Services/BackgroundWorkflowCancellationDispatcherTests.cs create mode 100644 test/unit/Elsa.Workflows.Runtime.UnitTests/Services/LocalWorkflowClientTests.cs diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs index fe1a4cc55..dcbc2ed55 100644 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs @@ -9,8 +9,9 @@ namespace Elsa.Api.Client.Resources.WorkflowDefinitions.Models; [PublicAPI] public class WorkflowDefinitionSummary : LinkedEntity { - public string DefinitionId { get; set; } = default!; - public string Name { get; set; } = default!; + public string DefinitionId { get; set; } = null!; + public string Name { get; set; } = null!; public string? Description { get; set; } - public string MaterializerName { get; set; } = default!; + public string MaterializerName { get; set; } = null!; + public bool IsMaterializerAvailable { get; set; } } \ No newline at end of file diff --git a/src/modules/Elsa.Caching/Contracts/ICacheManager.cs b/src/modules/Elsa.Caching/Contracts/ICacheManager.cs index 6e93fe8c8..1b75ffbed 100644 --- a/src/modules/Elsa.Caching/Contracts/ICacheManager.cs +++ b/src/modules/Elsa.Caching/Contracts/ICacheManager.cs @@ -25,8 +25,13 @@ public interface ICacheManager /// ValueTask TriggerTokenAsync(string key, CancellationToken cancellationToken = default); + /// + /// Finds an item from the cache, or creates it if it doesn't exist. + /// + Task FindOrCreateAsync(object key, Func> factory); + /// /// Gets an item from the cache, or creates it if it doesn't exist. /// - Task GetOrCreateAsync(object key, Func> factory); + Task GetOrCreateAsync(object key, Func> factory); } \ No newline at end of file diff --git a/src/modules/Elsa.Caching/Services/CacheManager.cs b/src/modules/Elsa.Caching/Services/CacheManager.cs index 6db501327..657a1c66d 100644 --- a/src/modules/Elsa.Caching/Services/CacheManager.cs +++ b/src/modules/Elsa.Caching/Services/CacheManager.cs @@ -24,8 +24,21 @@ public class CacheManager(IMemoryCache memoryCache, IChangeTokenSignaler changeT } /// - public async Task GetOrCreateAsync(object key, Func> factory) + public async Task FindOrCreateAsync(object key, Func> factory) { return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry)); } + + /// + /// Retrieves a cached item by the specified key or creates a new one using the provided factory function. + /// + /// The key used to identify the cached item. + /// A factory function that provides the value to be cached if it does not already exist. + /// The type of the item to retrieve or create. + /// The cached or newly created item. + /// Thrown if the factory function returns null. + public async Task GetOrCreateAsync(object key, Func> factory) + { + return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry)) ?? throw new InvalidOperationException($"Factory returned null for cache key: {key}."); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Common/Codecs/Zstd.cs b/src/modules/Elsa.Common/Codecs/Zstd.cs index 11a49c5d6..c8b46785a 100644 --- a/src/modules/Elsa.Common/Codecs/Zstd.cs +++ b/src/modules/Elsa.Common/Codecs/Zstd.cs @@ -15,7 +15,7 @@ public class Zstd : ICompressionCodec { var inputBytes = Encoding.UTF8.GetBytes(input); var span = inputBytes.AsSpan(); - var result = Iron.Compress(Codec.Zstd, span); + using var result = Iron.Compress(Codec.Zstd, span); var compressedBytes = result.AsSpan(); var compressedString = Convert.ToBase64String(compressedBytes); @@ -27,7 +27,7 @@ public class Zstd : ICompressionCodec { var inputBytes = Convert.FromBase64String(input); var span = inputBytes.AsSpan(); - var result = Iron.Decompress(Codec.Zstd, span); + using var result = Iron.Decompress(Codec.Zstd, span); var decompressedBytes = result.AsSpan(); var decompressedString = Encoding.UTF8.GetString(decompressedBytes); diff --git a/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs index b4d2e63ac..1667c45c5 100644 --- a/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs +++ b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs @@ -21,7 +21,7 @@ public class CachingHttpWorkflowLookupService( var tenantIdPrefix = !string.IsNullOrEmpty(tenantId) ? $"{tenantId}:" : string.Empty; var key = $"{tenantIdPrefix}http-workflow:{bookmarkHash}"; var cache = cacheManager.Cache; - return await cache.GetOrCreateAsync(key, async entry => + return await cache.FindOrCreateAsync(key, async entry => { var cachingOptions = cache.CachingOptions.Value; entry.SetSlidingExpiration(cachingOptions.CacheDuration); diff --git a/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs b/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs index 5c18cf948..67236da2c 100644 --- a/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs +++ b/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs @@ -1,4 +1,3 @@ -using Elsa.Dsl.ElsaScript.Compiler; using Elsa.Dsl.ElsaScript.Features; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs index 813ca17e7..93f333056 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs @@ -3,11 +3,12 @@ using Elsa.Common.Models; using Elsa.Workflows.Management; using Elsa.Workflows.Models; using JetBrains.Annotations; +using Microsoft.AspNetCore.Http; namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.GetByDefinitionId; [PublicAPI] -internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefinitionLinker linker) : ElsaEndpoint +internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefinitionLinker linker, IMaterializerRegistry materializerRegistry) : ElsaEndpoint { public override void Configure() { @@ -20,13 +21,20 @@ internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefini var versionOptions = request.VersionOptions != null ? VersionOptions.FromString(request.VersionOptions) : VersionOptions.Latest; var filter = WorkflowDefinitionHandle.ByDefinitionId(request.DefinitionId, versionOptions).ToFilter(); var definition = await store.FindAsync(filter, cancellationToken); - + if (definition == null) { await Send.NotFoundAsync(cancellationToken); return; } + if (!materializerRegistry.IsMaterializerAvailable(definition.MaterializerName)) + { + AddError($"The workflow materializer '{definition.MaterializerName}' is not available. The materializer may be disabled or not registered."); + await Send.ErrorsAsync(StatusCodes.Status422UnprocessableEntity, cancellationToken); + return; + } + var model = await linker.MapAsync(definition, cancellationToken); await Send.OkAsync(model, cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs b/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs index ebf28f35f..79a56cd68 100644 --- a/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs +++ b/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs @@ -6,4 +6,9 @@ namespace Elsa.Workflows.Api.Models; public class LinkedWorkflowDefinitionSummary : WorkflowDefinitionSummary { public Link[]? Links { get; set; } + + /// + /// Indicates whether the materializer for this workflow definition is currently available. + /// + public bool IsMaterializerAvailable { get; set; } } diff --git a/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs b/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs index da2a3351e..87ebc22ff 100644 --- a/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs +++ b/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs @@ -1,5 +1,6 @@ using Elsa.Models; using Elsa.Workflows.Api.Models; +using Elsa.Workflows.Management; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Mappers; using Elsa.Workflows.Management.Models; @@ -11,7 +12,8 @@ namespace Elsa.Workflows.Api.Services; /// public class StaticWorkflowDefinitionLinker( IOptions managementOptions, - WorkflowDefinitionMapper workflowDefinitionMapper) + WorkflowDefinitionMapper workflowDefinitionMapper, + IMaterializerRegistry materializerRegistry) : IWorkflowDefinitionLinker { /// @@ -51,7 +53,7 @@ public class StaticWorkflowDefinitionLinker( foreach (var item in list.Items) { - items.Add(new LinkedWorkflowDefinitionSummary + items.Add(new() { Links = GenerateLinksForSingleEntry(item.DefinitionId, item.IsReadonly), Id = item.Id, @@ -65,11 +67,12 @@ public class StaticWorkflowDefinitionLinker( ProviderName = item.ProviderName, MaterializerName = item.MaterializerName, CreatedAt = item.CreatedAt, - IsReadonly = item.IsReadonly + IsReadonly = item.IsReadonly, + IsMaterializerAvailable = materializerRegistry.IsMaterializerAvailable(item.MaterializerName) }); } - return new PagedListResponse + return new() { TotalCount = list.TotalCount, Items = items, @@ -126,13 +129,13 @@ public class StaticWorkflowDefinitionLinker( if (!managementOptions.Value.IsReadOnlyMode) { - linksList.Add(new Link($"/bulk-actions/delete/workflow-definitions/by-definition-id", "bulk-delete-by-definition-id", "POST")); - linksList.Add(new Link($"/bulk-actions/delete/workflow-definitions/by-id", "bulk-delete-by-id", "POST")); - linksList.Add(new Link($"/bulk-actions/publish/workflow-definitions/by-definition-ids", "bulk-publish", "POST")); - linksList.Add(new Link($"/bulk-actions/retract/workflow-definitions/by-definition-ids", "bulk-retract", "POST")); - linksList.Add(new Link($"/workflow-definitions/import", "import", "POST")); - linksList.Add(new Link($"/workflow-definitions/import-files", "import-files", "POST")); - linksList.Add(new Link($"/workflow-definitions", "create", "POST")); + linksList.Add(new($"/bulk-actions/delete/workflow-definitions/by-definition-id", "bulk-delete-by-definition-id", "POST")); + linksList.Add(new($"/bulk-actions/delete/workflow-definitions/by-id", "bulk-delete-by-id", "POST")); + linksList.Add(new($"/bulk-actions/publish/workflow-definitions/by-definition-ids", "bulk-publish", "POST")); + linksList.Add(new($"/bulk-actions/retract/workflow-definitions/by-definition-ids", "bulk-retract", "POST")); + linksList.Add(new($"/workflow-definitions/import", "import", "POST")); + linksList.Add(new($"/workflow-definitions/import-files", "import-files", "POST")); + linksList.Add(new($"/workflow-definitions", "create", "POST")); } return linksList.ToArray(); @@ -153,11 +156,11 @@ public class StaticWorkflowDefinitionLinker( if (!managementOptions.Value.IsReadOnlyMode && !definitionIsReadonly) { - links.Add(new Link($"/workflow-definitions/{definitionId}/publish", "publish", "POST")); - links.Add(new Link($"/workflow-definitions/{definitionId}/retract", "retract", "POST")); - links.Add(new Link($"/workflow-definitions/{definitionId}", "delete", "DELETE")); - links.Add(new Link($"/workflow-definitions/{definitionId}/import", "import", "PUT")); - links.Add(new Link($"/workflow-definitions/{definitionId}/update-references", "update-references", "POST")); + links.Add(new($"/workflow-definitions/{definitionId}/publish", "publish", "POST")); + links.Add(new($"/workflow-definitions/{definitionId}/retract", "retract", "POST")); + links.Add(new($"/workflow-definitions/{definitionId}", "delete", "DELETE")); + links.Add(new($"/workflow-definitions/{definitionId}/import", "import", "PUT")); + links.Add(new($"/workflow-definitions/{definitionId}/update-references", "update-references", "POST")); } return links.ToArray(); diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IMaterializerRegistry.cs b/src/modules/Elsa.Workflows.Management/Contracts/IMaterializerRegistry.cs new file mode 100644 index 000000000..b314e7cfc --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IMaterializerRegistry.cs @@ -0,0 +1,26 @@ +namespace Elsa.Workflows.Management; + +/// +/// Provides access to registered workflow materializers and their availability. +/// +public interface IMaterializerRegistry +{ + /// + /// Gets all registered materializers. + /// + IEnumerable GetMaterializers(); + + /// + /// Gets a materializer by name. + /// + /// The name of the materializer. + /// The materializer, or null if not found. + IWorkflowMaterializer? GetMaterializer(string name); + + /// + /// Checks if a materializer with the specified name is available. + /// + /// The name of the materializer. + /// True if the materializer is available, false otherwise. + bool IsMaterializerAvailable(string name); +} diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs index 149591861..0ffa08dd3 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs @@ -2,6 +2,7 @@ using Elsa.Common.Models; using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; using Elsa.Workflows.Models; namespace Elsa.Workflows.Management; @@ -60,4 +61,62 @@ public interface IWorkflowDefinitionService /// Looks for all s that match the specified . /// Task> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); + + /// + /// Attempts to retrieve a and its corresponding + /// based on the specified workflow definition ID and version options. + /// + /// The ID of the workflow definition to search for. + /// The versioning options that determine which version of the workflow definition to search for. + /// A token to monitor for cancellation requests. + /// + /// A containing the workflow graph and its associated definition, + /// or indicating that they do not exist. + /// + Task TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default); + + /// + /// Attempts to locate a based on the specified workflow definition version ID. + /// + /// The unique identifier of the workflow definition version. + /// A token to monitor for cancellation requests. + /// + /// A containing the workflow graph and its associated definition, + /// or indicating that they do not exist. + /// + Task TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default); + + /// + /// Attempts to retrieve the workflow graph corresponding to the specified . + /// + /// The handle identifying the workflow definition and its version. + /// A token to cancel the asynchronous operation. + /// + /// A containing the workflow graph and its associated definition, + /// or indicating that they do not exist. + /// + Task TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default); + + /// + /// Attempts to retrieve a and its associated + /// for the specified workflow definition using the provided filter criteria. + /// + /// The criteria used to filter and locate the workflow definition and graph. + /// A token to monitor for cancellation requests. + /// + /// A containing the workflow graph and its associated definition, + /// or indicating that they do not exist. + /// + Task TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); + + /// + /// Attempts to find and retrieve a collection of objects based on the specified . + /// + /// The filter specifying criteria to identify the workflow definitions to search for. + /// A token to observe while waiting for the task to complete. + /// + /// A collection of instances, each containing a workflow graph and its associated definition + /// for every matching workflow definition. If no matching workflow definitions are found, the collection is empty. + /// + Task> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Exceptions/WorkflowMaterializerNotFoundException.cs b/src/modules/Elsa.Workflows.Management/Exceptions/WorkflowMaterializerNotFoundException.cs new file mode 100644 index 000000000..118ec869d --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Exceptions/WorkflowMaterializerNotFoundException.cs @@ -0,0 +1,6 @@ +namespace Elsa.Workflows.Management.Exceptions; + +public class WorkflowMaterializerNotFoundException(string materializerName, string message = "Materializer not found. The materializer may be disabled or not registered") : Exception(message) +{ + public string MaterializerName { get; } = materializerName; +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index 4740bf1e5..cd3da8fec 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -255,6 +255,7 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module) .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs b/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs index 1d02f2970..8dafcd254 100644 --- a/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs +++ b/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs @@ -12,7 +12,7 @@ public class TimestampFilter /// /// Gets or sets the column to filter by. /// - public required string Column { get; set; } = default!; + public required string Column { get; set; } = null!; /// /// Gets or sets the operator to use for filtering. diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphFindResult.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphFindResult.cs new file mode 100644 index 000000000..d3df0d968 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphFindResult.cs @@ -0,0 +1,10 @@ +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Models; + +public record WorkflowGraphFindResult(WorkflowDefinition? WorkflowDefinition, WorkflowGraph? WorkflowGraph) +{ + public bool WorkflowDefinitionExists => WorkflowDefinition != null; + public bool WorkflowGraphExists => WorkflowGraph != null; +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs index 059457ba8..578c1fe84 100644 --- a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs @@ -1,6 +1,7 @@ using Elsa.Common.Models; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; using Elsa.Workflows.Models; using JetBrains.Annotations; using Microsoft.Extensions.Caching.Memory; @@ -11,12 +12,16 @@ namespace Elsa.Workflows.Management.Services; /// Decorates an with caching capabilities. /// [UsedImplicitly] -public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowDefinitionService +public class CachingWorkflowDefinitionService( + IWorkflowDefinitionService decoratedService, + IWorkflowDefinitionCacheManager cacheManager, + IWorkflowDefinitionStore workflowDefinitionStore, + IMaterializerRegistry materializerRegistry) : IWorkflowDefinitionService { /// - public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + public Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { - return await decoratedService.MaterializeWorkflowAsync(definition, cancellationToken); + return decoratedService.MaterializeWorkflowAsync(definition, cancellationToken); } /// @@ -24,7 +29,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat { var cacheKey = cacheManager.CreateWorkflowDefinitionVersionCacheKey(definitionId, versionOptions); - return await GetFromCacheAsync(cacheKey, + return await FindFromCacheAsync(cacheKey, () => decoratedService.FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken), x => x.DefinitionId); } @@ -33,7 +38,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat public async Task FindWorkflowDefinitionAsync(string definitionVersionId, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowDefinitionVersionCacheKey(definitionVersionId); - return await GetFromCacheAsync(cacheKey, + return await FindFromCacheAsync(cacheKey, () => decoratedService.FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken), x => x.DefinitionId); } @@ -41,8 +46,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat /// public Task FindWorkflowDefinitionAsync(WorkflowDefinitionHandle handle, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter { DefinitionHandle = handle }; - + var filter = handle.ToFilter(); return FindWorkflowDefinitionAsync(filter, cancellationToken); } @@ -50,7 +54,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat public async Task FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowDefinitionFilterCacheKey(filter); - return await GetFromCacheAsync(cacheKey, + return await FindFromCacheAsync(cacheKey, () => decoratedService.FindWorkflowDefinitionAsync(filter, cancellationToken), x => x.DefinitionId); } @@ -59,7 +63,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat public async Task FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionId, versionOptions); - return await GetFromCacheAsync(cacheKey, + return await FindFromCacheAsync(cacheKey, () => decoratedService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken), x => x.Workflow.Identity.DefinitionId); } @@ -68,7 +72,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat public async Task FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionVersionId); - return await GetFromCacheAsync(cacheKey, + return await FindFromCacheAsync(cacheKey, () => decoratedService.FindWorkflowGraphAsync(definitionVersionId, cancellationToken), x => x.Workflow.Identity.DefinitionId); } @@ -76,8 +80,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat /// public Task FindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter { DefinitionHandle = definitionHandle }; - + var filter = definitionHandle.ToFilter(); return FindWorkflowGraphAsync(filter, cancellationToken); } @@ -85,11 +88,12 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat public async Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowFilterCacheKey(filter); - return await GetFromCacheAsync(cacheKey, + return await FindFromCacheAsync(cacheKey, () => decoratedService.FindWorkflowGraphAsync(filter, cancellationToken), x => x.Workflow.Identity.DefinitionId); } + /// public async Task> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken); @@ -101,26 +105,102 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat cacheKey, async () => await MaterializeWorkflowAsync(workflowDefinition, cancellationToken), wf => wf.Workflow.Identity.DefinitionId); - workflowGraphs.Add(workflowGraph!); + workflowGraphs.Add(workflowGraph); } return workflowGraphs; } - private async Task GetFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc) + /// + public async Task TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) + { + var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionId, versionOptions); + var result = await GetFromCacheAsync(cacheKey, + () => decoratedService.TryFindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken), + x => x.WorkflowDefinition?.DefinitionId); + + return result; + } + + /// + public async Task TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default) + { + var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionVersionId); + var result = await GetFromCacheAsync(cacheKey, + () => decoratedService.TryFindWorkflowGraphAsync(definitionVersionId, cancellationToken), + x => x.WorkflowDefinition?.DefinitionId); + + return result; + } + + /// + public Task TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default) + { + var filter = definitionHandle.ToFilter(); + return TryFindWorkflowGraphAsync(filter, cancellationToken); + } + + /// + public async Task TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var cacheKey = cacheManager.CreateWorkflowFilterCacheKey(filter); + var result = await GetFromCacheAsync(cacheKey, + () => decoratedService.TryFindWorkflowGraphAsync(filter, cancellationToken), + x => x.WorkflowDefinition?.DefinitionId); + + return result; + } + + /// + public async Task> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken); + var results = new List(); + foreach (var workflowDefinition in workflowDefinitions) + { + if (!materializerRegistry.IsMaterializerAvailable(workflowDefinition.MaterializerName)) + { + var unavailableResult = new WorkflowGraphFindResult(workflowDefinition, null); + results.Add(unavailableResult); + continue; + } + + var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(workflowDefinition.Id); + var workflowGraph = await FindFromCacheAsync( + cacheKey, + async () => await MaterializeWorkflowAsync(workflowDefinition, cancellationToken), + wf => wf.Workflow.Identity.DefinitionId); + + var result = new WorkflowGraphFindResult(workflowDefinition, workflowGraph); + results.Add(result); + } + + return results; + } + + private async Task FindFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc) where T : class + { + return await GetFromCacheAsync( + cacheKey, + getObjectFunc, + obj => obj != null ? getChangeTokenKeyFunc(obj) : null); + } + + private async Task GetFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc) { var cache = cacheManager.Cache; return await cache.GetOrCreateAsync(cacheKey, async entry => { entry.SetAbsoluteExpiration(cache.CachingOptions.Value.CacheDuration); var obj = await getObjectFunc(); - - if (obj == null) - return default; - var changeTokenKeyInput = getChangeTokenKeyFunc(obj); - var changeTokenKey = cacheManager.CreateWorkflowDefinitionChangeTokenKey(changeTokenKeyInput); - entry.AddExpirationToken(cache.GetToken(changeTokenKey)); + + if (changeTokenKeyInput != null) + { + var changeTokenKey = cacheManager.CreateWorkflowDefinitionChangeTokenKey(changeTokenKeyInput); + entry.AddExpirationToken(cache.GetToken(changeTokenKey)); + } + return obj; }); } diff --git a/src/modules/Elsa.Workflows.Management/Services/MaterializerRegistry.cs b/src/modules/Elsa.Workflows.Management/Services/MaterializerRegistry.cs new file mode 100644 index 000000000..f7239e881 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/MaterializerRegistry.cs @@ -0,0 +1,22 @@ +namespace Elsa.Workflows.Management.Services; + +/// +public class MaterializerRegistry(Func> materializers) : IMaterializerRegistry +{ + private readonly Lazy> _materializers = new(() => materializers().ToArray()); + + /// + public IEnumerable GetMaterializers() => _materializers.Value; + + /// + public IWorkflowMaterializer? GetMaterializer(string name) + { + return _materializers.Value.FirstOrDefault(x => x.Name == name); + } + + /// + public bool IsMaterializerAvailable(string name) + { + return _materializers.Value.Any(x => x.Name == name); + } +} diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs index 104ea4e8c..a91107ee7 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs @@ -1,7 +1,10 @@ using Elsa.Common.Models; using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Exceptions; using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; using Elsa.Workflows.Models; +using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Management.Services; @@ -9,17 +12,17 @@ namespace Elsa.Workflows.Management.Services; public class WorkflowDefinitionService( IWorkflowDefinitionStore workflowDefinitionStore, IWorkflowGraphBuilder workflowGraphBuilder, - Func> materializers) + IMaterializerRegistry materializerRegistry, + ILogger logger) : IWorkflowDefinitionService { /// public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { - var workflowMaterializers = materializers(); - var materializer = workflowMaterializers.FirstOrDefault(x => x.Name == definition.MaterializerName); + var materializer = materializerRegistry.GetMaterializer(definition.MaterializerName); if (materializer == null) - throw new("Provider not found"); + throw new WorkflowMaterializerNotFoundException(definition.MaterializerName); var workflow = await materializer.MaterializeAsync(definition, cancellationToken); return await workflowGraphBuilder.BuildAsync(workflow, cancellationToken); @@ -56,44 +59,28 @@ public class WorkflowDefinitionService( public async Task FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { var definition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken); - - if (definition == null) - return null; - - return await MaterializeWorkflowAsync(definition, cancellationToken); + return await TryMaterializeWorkflowAsync(definition, cancellationToken); } /// public async Task FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default) { var definition = await FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken); - - if (definition == null) - return null; - - return await MaterializeWorkflowAsync(definition, cancellationToken); + return await TryMaterializeWorkflowAsync(definition, cancellationToken); } /// public async Task FindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default) { var definition = await FindWorkflowDefinitionAsync(definitionHandle, cancellationToken); - - if (definition == null) - return null; - - return await MaterializeWorkflowAsync(definition, cancellationToken); + return await TryMaterializeWorkflowAsync(definition, cancellationToken); } /// public async Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken); - - if (definition == null) - return null; - - return await MaterializeWorkflowAsync(definition, cancellationToken); + return await TryMaterializeWorkflowAsync(definition, cancellationToken); } /// @@ -109,4 +96,81 @@ public class WorkflowDefinitionService( return workflowGraphs; } + + public async Task TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) + { + var definition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken); + return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken); + } + + public async Task TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default) + { + var definition = await FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken); + return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken); + } + + public async Task TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default) + { + var definition = await FindWorkflowDefinitionAsync(definitionHandle, cancellationToken); + return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken); + } + + public async Task TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken); + return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken); + } + + public async Task> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken); + var results = new List(); + foreach (var workflowDefinition in workflowDefinitions) + { + var result = await MaterializeWorkflowGraphFindResultAsync(workflowDefinition, cancellationToken); + results.Add(result); + } + + return results; + } + + /// + /// Attempts to materialize a workflow graph from the given workflow definition if a suitable materializer is available. + /// + /// The workflow definition to materialize. Can be null. + /// A token to monitor for cancellation requests. + /// + /// A if materialization is successful; otherwise, null. + /// + private async Task TryMaterializeWorkflowAsync(WorkflowDefinition? definition, CancellationToken cancellationToken) + { + if (definition == null) + return null; + + if (materializerRegistry.IsMaterializerAvailable(definition.MaterializerName)) + return await MaterializeWorkflowAsync(definition, cancellationToken); + + logger.LogWarning("Materializer '{MaterializerName}' not found. The workflow definition will not be materialized.", definition.MaterializerName); + return null; + } + + /// + /// Attempts to materialize a workflow graph from the given workflow definition. + /// + /// The workflow definition to materialize the graph for. May be null. + /// A token to observe while waiting for the task to complete. + /// A result containing the materialized workflow graph and its corresponding workflow definition, if successfully materialized; otherwise, returns the definition with a null graph. + private async Task MaterializeWorkflowGraphFindResultAsync(WorkflowDefinition? definition, CancellationToken cancellationToken) + { + if (definition == null) + return new(null, null); + + if (materializerRegistry.IsMaterializerAvailable(definition.MaterializerName)) + { + var graph = await MaterializeWorkflowAsync(definition, cancellationToken); + return new(definition, graph); + } + + return new(definition, null); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs b/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs index 796a55755..22132cd2e 100644 --- a/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs +++ b/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs @@ -140,7 +140,7 @@ public class CachingWorkflowDefinitionStore(IWorkflowDefinitionStore decoratedSt var tenantId = tenantAccessor.Tenant?.Id; var tenantIdPrefix = !string.IsNullOrEmpty(tenantId) ? $"{tenantId}:" : string.Empty; var internalKey = $"{tenantIdPrefix}{typeof(T).Name}:{key}"; - return await cacheManager.GetOrCreateAsync(internalKey, async entry => + return await cacheManager.FindOrCreateAsync(internalKey, async entry => { var invalidationRequestToken = cacheManager.GetToken(CacheInvalidationTokenKey); entry.AddExpirationToken(invalidationRequestToken); diff --git a/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowDefinitionNotFoundException.cs b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowDefinitionNotFoundException.cs new file mode 100644 index 000000000..de7ac14d2 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowDefinitionNotFoundException.cs @@ -0,0 +1,8 @@ +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Runtime.Exceptions; + +public class WorkflowDefinitionNotFoundException(string message, WorkflowDefinitionHandle workflowDefinitionHandle) : Exception(message) +{ + public WorkflowDefinitionHandle WorkflowDefinitionHandle { get; } = workflowDefinitionHandle; +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs index 1d0f302fe..bee3b00a0 100644 --- a/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs +++ b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs @@ -2,6 +2,7 @@ using Elsa.Workflows.Models; namespace Elsa.Workflows.Runtime.Exceptions; +[Obsolete("Use WorkflowDefinitionNotFoundException instead.")] public class WorkflowGraphNotFoundException(string message, WorkflowDefinitionHandle workflowDefinitionHandle) : Exception(message) { public WorkflowDefinitionHandle WorkflowDefinitionHandle { get; } = workflowDefinitionHandle; diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs index a07d3823f..761f8396f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs @@ -1,21 +1,28 @@ -using Elsa.Mediator; -using Elsa.Mediator.Contracts; -using Elsa.Workflows.Runtime.Commands; -using Elsa.Workflows.Runtime.Requests; -using Elsa.Workflows.Runtime.Responses; - -namespace Elsa.Workflows.Runtime; - -/// -/// Dispatches workflow cancellation requests to a local background worker. -/// -public class BackgroundWorkflowCancellationDispatcher(ICommandSender commandSender) : IWorkflowCancellationDispatcher -{ - /// - public async Task DispatchAsync(DispatchCancelWorkflowRequest request, CancellationToken cancellationToken = default) - { - var command = new CancelWorkflowsCommand(request); - await commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken); - return new DispatchCancelWorkflowsResponse(); - } +using Elsa.Common.Multitenancy; +using Elsa.Mediator; +using Elsa.Mediator.Contracts; +using Elsa.Tenants.Mediator; +using Elsa.Workflows.Runtime.Commands; +using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Responses; + +namespace Elsa.Workflows.Runtime; + +/// +/// Dispatches workflow cancellation requests to a local background worker. +/// +public class BackgroundWorkflowCancellationDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IWorkflowCancellationDispatcher +{ + /// + public async Task DispatchAsync(DispatchCancelWorkflowRequest request, CancellationToken cancellationToken = default) + { + var command = new CancelWorkflowsCommand(request); + await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); + return new(); + } + + private IDictionary CreateHeaders() + { + return TenantHeaders.CreateHeaders(tenantAccessor.Tenant?.Id); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs b/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs index 3fc283d3e..75bcfb118 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs @@ -1,5 +1,6 @@ using Elsa.Workflows.Management; using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Exceptions; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Mappers; using Elsa.Workflows.Management.Options; @@ -207,8 +208,9 @@ public class LocalWorkflowClient( private async Task GetWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken) { - var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(definitionHandle, cancellationToken); - if (workflowGraph == null) throw new WorkflowGraphNotFoundException("Workflow graph not found.", definitionHandle); - return workflowGraph; + var result = await workflowDefinitionService.TryFindWorkflowGraphAsync(definitionHandle, cancellationToken); + if (!result.WorkflowDefinitionExists) throw new WorkflowDefinitionNotFoundException("Workflow definition not found.", definitionHandle); + if (!result.WorkflowGraphExists) throw new WorkflowMaterializerNotFoundException(result.WorkflowDefinition!.MaterializerName); + return result.WorkflowGraph!; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs index 5fed21c30..e84845dc9 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs @@ -77,7 +77,7 @@ public class CachingTriggerStore(ITriggerStore decoratedStore, ICacheManager cac var tenantId = tenantAccessor.Tenant?.Id; var tenantIdPrefix = !string.IsNullOrEmpty(tenantId) ? $"{tenantId}:" : string.Empty; var internalKey = $"{tenantIdPrefix}{typeof(T).Name}:{key}"; - return await cacheManager.GetOrCreateAsync(internalKey, async entry => + return await cacheManager.FindOrCreateAsync(internalKey, async entry => { var invalidationRequestToken = cacheManager.GetToken(CacheInvalidationTokenKey); entry.AddExpirationToken(invalidationRequestToken); diff --git a/test/unit/Elsa.Common.UnitTests/Codecs/ZstdTests.cs b/test/unit/Elsa.Common.UnitTests/Codecs/ZstdTests.cs new file mode 100644 index 000000000..3c6f7d908 --- /dev/null +++ b/test/unit/Elsa.Common.UnitTests/Codecs/ZstdTests.cs @@ -0,0 +1,67 @@ +using Elsa.Common.Codecs; + +namespace Elsa.Common.UnitTests.Codecs; + +public class ZstdTests +{ + private readonly Zstd _codec = new(); + + [Fact] + public async Task CompressAsync_WithSimpleString_ReturnsCompressedString() + { + // Arrange + var input = "Hello, World!"; + + // Act + var result = await _codec.CompressAsync(input); + + // Assert + Assert.NotNull(result); + Assert.NotEmpty(result); + Assert.NotEqual(input, result); + } + + [Theory] + [InlineData("Hello, World!")] + [InlineData("")] + [InlineData("Hello! 你好! مرحبا! Здравствуйте! 🎉🎊")] + [InlineData("{\"name\":\"John Doe\",\"age\":30,\"city\":\"New York\",\"items\":[1,2,3,4,5]}")] + public async Task CompressDecompress_RoundTrip_PreservesOriginalData(string original) + { + // Act + var compressed = await _codec.CompressAsync(original); + var decompressed = await _codec.DecompressAsync(compressed); + + // Assert + Assert.Equal(original, decompressed); + } + + [Fact] + public async Task CompressDecompress_WithLargeString_WorksCorrectly() + { + // Arrange + var original = string.Join("", Enumerable.Repeat("This is a test string that will be compressed. ", 1000)); + + // Act + var compressed = await _codec.CompressAsync(original); + var decompressed = await _codec.DecompressAsync(compressed); + + // Assert + Assert.Equal(original, decompressed); + Assert.True(compressed.Length < original.Length, "Compressed string should be smaller than original"); + } + + [Fact] + public async Task CompressAsync_MultipleCallsWithSameInput_ProducesConsistentResults() + { + // Arrange + var input = "Test string for consistency"; + + // Act + var result1 = await _codec.CompressAsync(input); + var result2 = await _codec.CompressAsync(input); + + // Assert + Assert.Equal(result1, result2); + } +} diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Helpers/TestHelpers.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Helpers/TestHelpers.cs new file mode 100644 index 000000000..04b70a82b --- /dev/null +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Helpers/TestHelpers.cs @@ -0,0 +1,40 @@ +using Elsa.Workflows.Activities; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.UnitTests.Helpers; + +/// +/// Provides helper methods for creating test data in Workflow Management unit tests. +/// +public static class TestHelpers +{ + /// + /// Creates a workflow definition with the specified parameters. + /// + public static WorkflowDefinition CreateWorkflowDefinition(string definitionId, string materializerName) + { + return new WorkflowDefinition + { + DefinitionId = definitionId, + MaterializerName = materializerName, + Version = 1 + }; + } + + /// + /// Creates a workflow graph from a workflow, ensuring proper initialization. + /// + public static WorkflowGraph CreateWorkflowGraph(Workflow workflow) + { + // Ensure the workflow has a proper root activity with an ID + var rootActivity = new Sequence { Id = "Root" }; + workflow.Root = rootActivity; + + // Create an ActivityNode with the root activity + var rootNode = new ActivityNode(rootActivity, "Root"); + + // Create and return the WorkflowGraph + return new WorkflowGraph(workflow, rootNode, new[] { rootNode }); + } +} diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Services/CachingWorkflowDefinitionServiceTests.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Services/CachingWorkflowDefinitionServiceTests.cs new file mode 100644 index 000000000..9204667f7 --- /dev/null +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Services/CachingWorkflowDefinitionServiceTests.cs @@ -0,0 +1,487 @@ +using Elsa.Caching; +using Elsa.Caching.Options; +using Elsa.Common.Models; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Management.Services; +using Elsa.Workflows.Management.UnitTests.Helpers; +using Elsa.Workflows.Models; +using Microsoft.Extensions.Caching.Memory; +using NSubstitute; + +namespace Elsa.Workflows.Management.UnitTests.Services; + +public class CachingWorkflowDefinitionServiceTests +{ + private readonly IWorkflowDefinitionService _decoratedService = Substitute.For(); + private readonly IWorkflowDefinitionCacheManager _cacheManager = Substitute.For(); + private readonly IWorkflowDefinitionStore _workflowDefinitionStore = Substitute.For(); + private readonly IMaterializerRegistry _materializerRegistry = Substitute.For(); + private readonly ICacheManager _cache = Substitute.For(); + + public CachingWorkflowDefinitionServiceTests() + { + SetupCacheManager(); + } + + [Fact] + public async Task MaterializeWorkflowAsync_DelegatesToDecoratedService() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var (_, workflowGraph) = CreateWorkflowAndGraph("def-1"); + + _decoratedService.MaterializeWorkflowAsync(definition, Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.MaterializeWorkflowAsync(definition); + + // Assert + Assert.Same(workflowGraph, result); + await _decoratedService.Received(1).MaterializeWorkflowAsync(definition, Arg.Any()); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByDefinitionIdAndVersionOptions_CreatesCacheKey() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var cacheKey = "cache-key-1"; + + _cacheManager.CreateWorkflowDefinitionVersionCacheKey("def-1", VersionOptions.Published).Returns(cacheKey); + _decoratedService.FindWorkflowDefinitionAsync("def-1", VersionOptions.Published, Arg.Any()) + .Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync("def-1", VersionOptions.Published); + + // Assert + Assert.Same(definition, result); + _cacheManager.Received(1).CreateWorkflowDefinitionVersionCacheKey("def-1", VersionOptions.Published); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByDefinitionVersionId_CreatesCacheKey() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + definition.Id = "version-id-1"; + var cacheKey = "cache-key-version-1"; + + _cacheManager.CreateWorkflowDefinitionVersionCacheKey("version-id-1").Returns(cacheKey); + _decoratedService.FindWorkflowDefinitionAsync("version-id-1", Arg.Any()).Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync("version-id-1"); + + // Assert + Assert.Same(definition, result); + _cacheManager.Received(1).CreateWorkflowDefinitionVersionCacheKey("version-id-1"); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByHandle_ConvertsTFilterAndDelegates() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var handle = WorkflowDefinitionHandle.ByDefinitionId("def-1", VersionOptions.Latest); + var cacheKey = "cache-key-filter"; + + _cacheManager.CreateWorkflowDefinitionFilterCacheKey(Arg.Any()).Returns(cacheKey); + _decoratedService.FindWorkflowDefinitionAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync(handle); + + // Assert + Assert.Same(definition, result); + await _decoratedService.Received(1).FindWorkflowDefinitionAsync(Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByFilter_CreatesCacheKey() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var filter = new WorkflowDefinitionFilter { DefinitionId = "def-1" }; + var cacheKey = "cache-key-filter"; + + _cacheManager.CreateWorkflowDefinitionFilterCacheKey(filter).Returns(cacheKey); + _decoratedService.FindWorkflowDefinitionAsync(filter, Arg.Any()).Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync(filter); + + // Assert + Assert.Same(definition, result); + _cacheManager.Received(1).CreateWorkflowDefinitionFilterCacheKey(filter); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_CreatesCacheKey() + { + // Arrange + var (_, workflowGraph) = CreateWorkflowAndGraph("def-1"); + var cacheKey = "cache-key-graph"; + + _cacheManager.CreateWorkflowVersionCacheKey("def-1", VersionOptions.Published).Returns(cacheKey); + _decoratedService.FindWorkflowGraphAsync("def-1", VersionOptions.Published, Arg.Any()) + .Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.Same(workflowGraph, result); + _cacheManager.Received(1).CreateWorkflowVersionCacheKey("def-1", VersionOptions.Published); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByDefinitionVersionId_CreatesCacheKey() + { + // Arrange + var (_, workflowGraph) = CreateWorkflowAndGraph("def-1"); + var cacheKey = "cache-key-graph-version"; + + _cacheManager.CreateWorkflowVersionCacheKey("version-id-1").Returns(cacheKey); + _decoratedService.FindWorkflowGraphAsync("version-id-1", Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync("version-id-1"); + + // Assert + Assert.Same(workflowGraph, result); + _cacheManager.Received(1).CreateWorkflowVersionCacheKey("version-id-1"); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByHandle_ConvertsTFilterAndDelegates() + { + // Arrange + var (_, workflowGraph) = CreateWorkflowAndGraph("def-1"); + var handle = WorkflowDefinitionHandle.ByDefinitionId("def-1", VersionOptions.Latest); + var cacheKey = "cache-key-filter"; + + _cacheManager.CreateWorkflowFilterCacheKey(Arg.Any()).Returns(cacheKey); + _decoratedService.FindWorkflowGraphAsync(Arg.Any(), Arg.Any()) + .Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync(handle); + + // Assert + Assert.Same(workflowGraph, result); + await _decoratedService.Received(1).FindWorkflowGraphAsync(Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByFilter_CreatesCacheKey() + { + // Arrange + var (_, workflowGraph) = CreateWorkflowAndGraph("def-1"); + var filter = new WorkflowDefinitionFilter { DefinitionId = "def-1" }; + var cacheKey = "cache-key-filter"; + + _cacheManager.CreateWorkflowFilterCacheKey(filter).Returns(cacheKey); + _decoratedService.FindWorkflowGraphAsync(filter, Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync(filter); + + // Assert + Assert.Same(workflowGraph, result); + _cacheManager.Received(1).CreateWorkflowFilterCacheKey(filter); + } + + [Fact] + public async Task FindWorkflowGraphsAsync_WithMultipleDefinitions_CachesEachGraph() + { + // Arrange + var definition1 = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + definition1.Id = "id-1"; + var definition2 = TestHelpers.CreateWorkflowDefinition("def-2", "materializer"); + definition2.Id = "id-2"; + var definitions = new[] { definition1, definition2 }; + + var (_, graph1) = CreateWorkflowAndGraph("def-1"); + var (_, graph2) = CreateWorkflowAndGraph("def-2"); + + var filter = new WorkflowDefinitionFilter(); + _workflowDefinitionStore.FindManyAsync(filter, Arg.Any()).Returns(definitions); + + _cacheManager.CreateWorkflowVersionCacheKey("id-1").Returns("cache-key-1"); + _cacheManager.CreateWorkflowVersionCacheKey("id-2").Returns("cache-key-2"); + + _decoratedService.MaterializeWorkflowAsync(definition1, Arg.Any()).Returns(graph1); + _decoratedService.MaterializeWorkflowAsync(definition2, Arg.Any()).Returns(graph2); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphsAsync(filter); + + // Assert + var graphs = result.ToList(); + Assert.Equal(2, graphs.Count); + Assert.Contains(graph1, graphs); + Assert.Contains(graph2, graphs); + _cacheManager.Received(1).CreateWorkflowVersionCacheKey("id-1"); + _cacheManager.Received(1).CreateWorkflowVersionCacheKey("id-2"); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_CreatesCacheKey() + { + // Arrange + var findResult = CreateWorkflowGraphFindResult("def-1"); + var cacheKey = "cache-key-try"; + + _cacheManager.CreateWorkflowVersionCacheKey("def-1", VersionOptions.Published).Returns(cacheKey); + _decoratedService.TryFindWorkflowGraphAsync("def-1", VersionOptions.Published, Arg.Any()) + .Returns(findResult); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.Same(findResult, result); + _cacheManager.Received(1).CreateWorkflowVersionCacheKey("def-1", VersionOptions.Published); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByDefinitionVersionId_CreatesCacheKey() + { + // Arrange + var findResult = CreateWorkflowGraphFindResult("def-1"); + var cacheKey = "cache-key-try-version"; + + _cacheManager.CreateWorkflowVersionCacheKey("version-id-1").Returns(cacheKey); + _decoratedService.TryFindWorkflowGraphAsync("version-id-1", Arg.Any()).Returns(findResult); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync("version-id-1"); + + // Assert + Assert.Same(findResult, result); + _cacheManager.Received(1).CreateWorkflowVersionCacheKey("version-id-1"); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByHandle_ConvertsTFilterAndDelegates() + { + // Arrange + var findResult = CreateWorkflowGraphFindResult("def-1"); + var handle = WorkflowDefinitionHandle.ByDefinitionId("def-1", VersionOptions.Latest); + var cacheKey = "cache-key-filter"; + + _cacheManager.CreateWorkflowFilterCacheKey(Arg.Any()).Returns(cacheKey); + _decoratedService.TryFindWorkflowGraphAsync(Arg.Any(), Arg.Any()) + .Returns(findResult); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync(handle); + + // Assert + Assert.Same(findResult, result); + await _decoratedService.Received(1).TryFindWorkflowGraphAsync(Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByFilter_CreatesCacheKey() + { + // Arrange + var findResult = CreateWorkflowGraphFindResult("def-1"); + var filter = new WorkflowDefinitionFilter { DefinitionId = "def-1" }; + var cacheKey = "cache-key-filter"; + + _cacheManager.CreateWorkflowFilterCacheKey(filter).Returns(cacheKey); + _decoratedService.TryFindWorkflowGraphAsync(filter, Arg.Any()).Returns(findResult); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync(filter); + + // Assert + Assert.Same(findResult, result); + _cacheManager.Received(1).CreateWorkflowFilterCacheKey(filter); + } + + [Fact] + public async Task TryFindWorkflowGraphsAsync_WithAvailableMaterializer_CachesGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + definition.Id = "id-1"; + var definitions = new[] { definition }; + + var (_, graph) = CreateWorkflowAndGraph("def-1"); + + var filter = new WorkflowDefinitionFilter(); + _workflowDefinitionStore.FindManyAsync(filter, Arg.Any()).Returns(definitions); + + _cacheManager.CreateWorkflowVersionCacheKey("id-1").Returns("cache-key-1"); + _materializerRegistry.IsMaterializerAvailable("materializer").Returns(true); + _decoratedService.MaterializeWorkflowAsync(definition, Arg.Any()).Returns(graph); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphsAsync(filter); + + // Assert + var results = result.ToList(); + Assert.Single(results); + Assert.Same(definition, results[0].WorkflowDefinition); + Assert.Same(graph, results[0].WorkflowGraph); + } + + [Fact] + public async Task TryFindWorkflowGraphsAsync_WithUnavailableMaterializer_ReturnsNullGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "unavailable-materializer"); + definition.Id = "id-1"; + var definitions = new[] { definition }; + + var filter = new WorkflowDefinitionFilter(); + _workflowDefinitionStore.FindManyAsync(filter, Arg.Any()).Returns(definitions); + + _cacheManager.CreateWorkflowVersionCacheKey("id-1").Returns("cache-key-1"); + _materializerRegistry.IsMaterializerAvailable("unavailable-materializer").Returns(false); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphsAsync(filter); + + // Assert + var results = result.ToList(); + Assert.Single(results); + Assert.Same(definition, results[0].WorkflowDefinition); + Assert.Null(results[0].WorkflowGraph); + await _decoratedService.DidNotReceive().MaterializeWorkflowAsync(Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task TryFindWorkflowGraphsAsync_WithMixedMaterializers_HandlesEachCorrectly() + { + // Arrange + var definition1 = TestHelpers.CreateWorkflowDefinition("def-1", "available-materializer"); + definition1.Id = "id-1"; + var definition2 = TestHelpers.CreateWorkflowDefinition("def-2", "unavailable-materializer"); + definition2.Id = "id-2"; + var definitions = new[] { definition1, definition2 }; + + var (_, graph1) = CreateWorkflowAndGraph("def-1"); + + var filter = new WorkflowDefinitionFilter(); + _workflowDefinitionStore.FindManyAsync(filter, Arg.Any()).Returns(definitions); + + _cacheManager.CreateWorkflowVersionCacheKey("id-1").Returns("cache-key-1"); + _cacheManager.CreateWorkflowVersionCacheKey("id-2").Returns("cache-key-2"); + _materializerRegistry.IsMaterializerAvailable("available-materializer").Returns(true); + _materializerRegistry.IsMaterializerAvailable("unavailable-materializer").Returns(false); + _decoratedService.MaterializeWorkflowAsync(definition1, Arg.Any()).Returns(graph1); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphsAsync(filter); + + // Assert + var results = result.ToList(); + Assert.Equal(2, results.Count); + + var result1 = results.First(r => r.WorkflowDefinition?.DefinitionId == "def-1"); + Assert.Same(definition1, result1.WorkflowDefinition); + Assert.Same(graph1, result1.WorkflowGraph); + + var result2 = results.First(r => r.WorkflowDefinition?.DefinitionId == "def-2"); + Assert.Same(definition2, result2.WorkflowDefinition); + Assert.Null(result2.WorkflowGraph); + } + + private CachingWorkflowDefinitionService CreateService() + { + return new( + _decoratedService, + _cacheManager, + _workflowDefinitionStore, + _materializerRegistry); + } + + /// + /// Sets up the cache manager with proper mock behavior for all cache operations. + /// + private void SetupCacheManager() + { + _cacheManager.Cache.Returns(_cache); + + var cachingOptions = Microsoft.Extensions.Options.Options.Create(new CachingOptions { CacheDuration = TimeSpan.FromMinutes(10) }); + _cache.CachingOptions.Returns(cachingOptions); + + // Setup cache to call factory functions for all types + _cache.GetOrCreateAsync(Arg.Any(), Arg.Any>>()) + .Returns(async callInfo => await callInfo.Arg>>()(Substitute.For())); + + _cache.GetOrCreateAsync(Arg.Any(), Arg.Any>>()) + .Returns(async callInfo => await callInfo.Arg>>()(Substitute.For())); + + _cache.GetOrCreateAsync(Arg.Any(), Arg.Any>>()) + .Returns(async callInfo => await callInfo.Arg>>()(Substitute.For())); + + _cache.FindOrCreateAsync(Arg.Any(), Arg.Any>>()) + .Returns(async callInfo => await callInfo.Arg>>()(Substitute.For())); + + _cache.FindOrCreateAsync(Arg.Any(), Arg.Any>>()) + .Returns(async callInfo => await callInfo.Arg>>()(Substitute.For())); + } + + /// + /// Creates a workflow and workflow graph for testing. + /// + private (Workflow, WorkflowGraph) CreateWorkflowAndGraph(string definitionId = "def-1", int version = 1) + { + var workflow = new Workflow { Identity = new(definitionId, version, $"v{version}") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + return (workflow, workflowGraph); + } + + /// + /// Creates a workflow graph find result for testing. + /// + private WorkflowGraphFindResult CreateWorkflowGraphFindResult(string definitionId = "def-1", string materializerName = "materializer") + { + var definition = TestHelpers.CreateWorkflowDefinition(definitionId, materializerName); + var (_, workflowGraph) = CreateWorkflowAndGraph(definitionId); + return new(definition, workflowGraph); + } +} diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Services/MaterializerRegistryTests.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Services/MaterializerRegistryTests.cs new file mode 100644 index 000000000..cc38bc3ac --- /dev/null +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Services/MaterializerRegistryTests.cs @@ -0,0 +1,179 @@ +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Services; +using NSubstitute; + +namespace Elsa.Workflows.Management.UnitTests.Services; + +public class MaterializerRegistryTests +{ + [Fact] + public void GetMaterializers_Should_Return_All_Materializers() + { + // Arrange + var materializer1 = CreateMaterializer("materializer1"); + var materializer2 = CreateMaterializer("materializer2"); + var registry = CreateRegistry(materializer1, materializer2); + + // Act + var result = registry.GetMaterializers().ToList(); + + // Assert + Assert.Equal(2, result.Count); + Assert.Contains(result, m => m.Name == "materializer1"); + Assert.Contains(result, m => m.Name == "materializer2"); + } + + [Fact] + public void GetMaterializer_Should_Return_Materializer_By_Name() + { + // Arrange + var materializer1 = CreateMaterializer("test-materializer"); + var materializer2 = CreateMaterializer("other-materializer"); + var registry = CreateRegistry(materializer1, materializer2); + + // Act + var result = registry.GetMaterializer("test-materializer"); + + // Assert + Assert.NotNull(result); + Assert.Equal("test-materializer", result.Name); + } + + [Theory] + [InlineData("non-existent")] + [InlineData("")] + [InlineData("wrong-name")] + public void GetMaterializer_Should_Return_Null_When_Not_Found(string name) + { + // Arrange + var materializer = CreateMaterializer("existing-materializer"); + var registry = CreateRegistry(materializer); + + // Act + var result = registry.GetMaterializer(name); + + // Assert + Assert.Null(result); + } + + [Fact] + public void IsMaterializerAvailable_Should_Return_True_When_Materializer_Exists() + { + // Arrange + var materializer = CreateMaterializer("available-materializer"); + var registry = CreateRegistry(materializer); + + // Act + var result = registry.IsMaterializerAvailable("available-materializer"); + + // Assert + Assert.True(result); + } + + [Theory] + [InlineData("not-available")] + [InlineData("")] + [InlineData("wrong-name")] + public void IsMaterializerAvailable_Should_Return_False_When_Materializer_Does_Not_Exist(string name) + { + // Arrange + var materializer = CreateMaterializer("existing-materializer"); + var registry = CreateRegistry(materializer); + + // Act + var result = registry.IsMaterializerAvailable(name); + + // Assert + Assert.False(result); + } + + [Fact] + public void Should_Handle_Empty_Materializer_Collection() + { + // Arrange + var registry = CreateRegistry(); + + // Act + var allMaterializers = registry.GetMaterializers().ToList(); + var foundMaterializer = registry.GetMaterializer("any-name"); + var isAvailable = registry.IsMaterializerAvailable("any-name"); + + // Assert + Assert.Empty(allMaterializers); + Assert.Null(foundMaterializer); + Assert.False(isAvailable); + } + + [Fact] + public void Should_Cache_Materializers_On_First_Access() + { + // Arrange + var callCount = 0; + IEnumerable MaterializersFactory() + { + callCount++; + return new[] { CreateMaterializer("test") }; + } + + var registry = new MaterializerRegistry(MaterializersFactory); + + // Act + var first = registry.GetMaterializers().ToList(); + var second = registry.GetMaterializers().ToList(); + var third = registry.GetMaterializer("test"); + var fourth = registry.IsMaterializerAvailable("test"); + + // Assert + Assert.Single(first); + Assert.Single(second); + Assert.NotNull(third); + Assert.True(fourth); + Assert.Equal(1, callCount); // Factory should only be called once + } + + [Fact] + public void GetMaterializers_Should_Allow_Multiple_Enumerations() + { + // Arrange + var materializer1 = CreateMaterializer("mat1"); + var materializer2 = CreateMaterializer("mat2"); + var registry = CreateRegistry(materializer1, materializer2); + + // Act + var first = registry.GetMaterializers().ToList(); + var second = registry.GetMaterializers().ToList(); + + // Assert + Assert.Equal(2, first.Count); + Assert.Equal(2, second.Count); + Assert.Equal(first.Count, second.Count); + } + + [Fact] + public void GetMaterializer_Should_Return_First_When_Multiple_Materializers_With_Same_Name() + { + // Arrange - FirstOrDefault returns the first match + var materializer1 = CreateMaterializer("duplicate-name"); + var materializer2 = CreateMaterializer("duplicate-name"); + var registry = CreateRegistry(materializer1, materializer2); + + // Act + var result = registry.GetMaterializer("duplicate-name"); + + // Assert + Assert.NotNull(result); + Assert.Same(materializer1, result); // FirstOrDefault returns the first match + } + + private static IWorkflowMaterializer CreateMaterializer(string name) + { + var materializer = Substitute.For(); + materializer.Name.Returns(name); + return materializer; + } + + private static MaterializerRegistry CreateRegistry(params IWorkflowMaterializer[] materializers) + { + return new(() => materializers); + } +} diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionServiceTests.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionServiceTests.cs new file mode 100644 index 000000000..4a6146c58 --- /dev/null +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionServiceTests.cs @@ -0,0 +1,524 @@ +using Elsa.Common.Models; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Exceptions; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Services; +using Elsa.Workflows.Management.UnitTests.Helpers; +using Elsa.Workflows.Models; +using Microsoft.Extensions.Logging; +using NSubstitute; + +namespace Elsa.Workflows.Management.UnitTests.Services; + +public class WorkflowDefinitionServiceTests +{ + private readonly IWorkflowDefinitionStore _workflowDefinitionStore = Substitute.For(); + private readonly IWorkflowGraphBuilder _workflowGraphBuilder = Substitute.For(); + private readonly IMaterializerRegistry _materializerRegistry = Substitute.For(); + private readonly ILogger _logger = Substitute.For>(); + + [Fact] + public async Task MaterializeWorkflowAsync_WithValidMaterializer_ReturnsMaterializedGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "test-materializer"); + var (_, workflowGraph) = SetupMaterializerAndGraphBuilder(definition); + + var service = CreateService(); + + // Act + var result = await service.MaterializeWorkflowAsync(definition); + + // Assert + Assert.Same(workflowGraph, result); + } + + [Fact] + public async Task MaterializeWorkflowAsync_WithMissingMaterializer_ThrowsException() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "missing-materializer"); + _materializerRegistry.GetMaterializer("missing-materializer").Returns((IWorkflowMaterializer?)null); + + var service = CreateService(); + + // Act & Assert + await Assert.ThrowsAsync(() => service.MaterializeWorkflowAsync(definition)); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByDefinitionIdAndVersionOptions_ReturnsDefinition() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync("def-1", VersionOptions.Published); + + // Assert + Assert.Same(definition, result); + await _workflowDefinitionStore.Received(1).FindAsync( + Arg.Any(), + Arg.Any()); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByDefinitionVersionId_ReturnsDefinition() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + definition.Id = "version-id-1"; + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync("version-id-1"); + + // Assert + Assert.Same(definition, result); + await _workflowDefinitionStore.Received(1).FindAsync( + Arg.Is(f => f.Id == "version-id-1"), + Arg.Any()); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByHandle_ReturnsDefinition() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var handle = WorkflowDefinitionHandle.ByDefinitionId("def-1", VersionOptions.Latest); + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync(handle); + + // Assert + Assert.Same(definition, result); + await _workflowDefinitionStore.Received(1).FindAsync(Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task FindWorkflowDefinitionAsync_ByFilter_ReturnsDefinition() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var filter = new WorkflowDefinitionFilter { DefinitionId = "def-1" }; + _workflowDefinitionStore.FindAsync(filter, Arg.Any()).Returns(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowDefinitionAsync(filter); + + // Assert + Assert.Same(definition, result); + await _workflowDefinitionStore.Received(1).FindAsync(filter, Arg.Any()); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_WithAvailableMaterializer_ReturnsGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var workflowGraph = SetupCompleteWorkflowGraphMocking(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.Same(workflowGraph, result); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_WithUnavailableMaterializer_ReturnsNull() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + SetupCompleteWorkflowGraphMocking(definition, materializerAvailable: false); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.Null(result); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_WhenDefinitionNotFound_ReturnsNull() + { + // Arrange + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns((WorkflowDefinition?)null); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.Null(result); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByDefinitionVersionId_ReturnsGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + definition.Id = "version-id-1"; + var workflowGraph = SetupCompleteWorkflowGraphMocking(definition); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync("version-id-1"); + + // Assert + Assert.Same(workflowGraph, result); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByHandle_ReturnsGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var handle = WorkflowDefinitionHandle.ByDefinitionId("def-1", VersionOptions.Latest); + var workflow = new Workflow { Identity = new("def-1", 1, "v1") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + _materializerRegistry.IsMaterializerAvailable("materializer").Returns(true); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition, Arg.Any()).Returns(workflow); + _materializerRegistry.GetMaterializer("materializer").Returns(materializer); + _workflowGraphBuilder.BuildAsync(workflow, Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync(handle); + + // Assert + Assert.Same(workflowGraph, result); + } + + [Fact] + public async Task FindWorkflowGraphAsync_ByFilter_ReturnsGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var filter = new WorkflowDefinitionFilter { DefinitionId = "def-1" }; + var workflow = new Workflow { Identity = new("def-1", 1, "v1") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + + _workflowDefinitionStore.FindAsync(filter, Arg.Any()).Returns(definition); + _materializerRegistry.IsMaterializerAvailable("materializer").Returns(true); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition, Arg.Any()).Returns(workflow); + _materializerRegistry.GetMaterializer("materializer").Returns(materializer); + _workflowGraphBuilder.BuildAsync(workflow, Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphAsync(filter); + + // Assert + Assert.Same(workflowGraph, result); + } + + [Fact] + public async Task FindWorkflowGraphsAsync_WithMultipleDefinitions_ReturnsAllGraphs() + { + // Arrange + var definition1 = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var definition2 = TestHelpers.CreateWorkflowDefinition("def-2", "materializer"); + var definitions = new[] { definition1, definition2 }; + + var workflow1 = new Workflow { Identity = new("def-1", 1, "v1") }; + var workflow2 = new Workflow { Identity = new("def-2", 1, "v1") }; + var graph1 = TestHelpers.CreateWorkflowGraph(workflow1); + var graph2 = TestHelpers.CreateWorkflowGraph(workflow2); + + var filter = new WorkflowDefinitionFilter(); + _workflowDefinitionStore.FindManyAsync(filter, Arg.Any()).Returns(definitions); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition1, Arg.Any()).Returns(workflow1); + materializer.MaterializeAsync(definition2, Arg.Any()).Returns(workflow2); + _materializerRegistry.GetMaterializer("materializer").Returns(materializer); + + _workflowGraphBuilder.BuildAsync(workflow1, Arg.Any()).Returns(graph1); + _workflowGraphBuilder.BuildAsync(workflow2, Arg.Any()).Returns(graph2); + + var service = CreateService(); + + // Act + var result = await service.FindWorkflowGraphsAsync(filter); + + // Assert + var graphs = result.ToList(); + Assert.Equal(2, graphs.Count); + Assert.Contains(graph1, graphs); + Assert.Contains(graph2, graphs); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_WithAvailableMaterializer_ReturnsSuccessResult() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var workflow = new Workflow { Identity = new("def-1", 1, "v1") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + _materializerRegistry.IsMaterializerAvailable("materializer").Returns(true); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition, Arg.Any()).Returns(workflow); + _materializerRegistry.GetMaterializer("materializer").Returns(materializer); + _workflowGraphBuilder.BuildAsync(workflow, Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.NotNull(result); + Assert.Same(definition, result.WorkflowDefinition); + Assert.Same(workflowGraph, result.WorkflowGraph); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_WithUnavailableMaterializer_ReturnsResultWithNullGraph() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + SetupCompleteWorkflowGraphMocking(definition, materializerAvailable: false); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.NotNull(result); + Assert.Same(definition, result.WorkflowDefinition); + Assert.Null(result.WorkflowGraph); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByDefinitionIdAndVersionOptions_WhenDefinitionNotFound_ReturnsEmptyResult() + { + // Arrange + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns((WorkflowDefinition?)null); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync("def-1", VersionOptions.Published); + + // Assert + Assert.NotNull(result); + Assert.Null(result.WorkflowDefinition); + Assert.Null(result.WorkflowGraph); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByDefinitionVersionId_ReturnsResult() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + definition.Id = "version-id-1"; + var workflowGraph = SetupCompleteWorkflowGraphMocking(definition); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync("version-id-1"); + + // Assert + Assert.NotNull(result); + Assert.Same(definition, result.WorkflowDefinition); + Assert.Same(workflowGraph, result.WorkflowGraph); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByHandle_ReturnsResult() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var handle = WorkflowDefinitionHandle.ByDefinitionId("def-1", VersionOptions.Latest); + var workflow = new Workflow { Identity = new("def-1", 1, "v1") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + _materializerRegistry.IsMaterializerAvailable("materializer").Returns(true); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition, Arg.Any()).Returns(workflow); + _materializerRegistry.GetMaterializer("materializer").Returns(materializer); + _workflowGraphBuilder.BuildAsync(workflow, Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync(handle); + + // Assert + Assert.NotNull(result); + Assert.Same(definition, result.WorkflowDefinition); + Assert.Same(workflowGraph, result.WorkflowGraph); + } + + [Fact] + public async Task TryFindWorkflowGraphAsync_ByFilter_ReturnsResult() + { + // Arrange + var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var filter = new WorkflowDefinitionFilter { DefinitionId = "def-1" }; + var workflow = new Workflow { Identity = new("def-1", 1, "v1") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + + _workflowDefinitionStore.FindAsync(filter, Arg.Any()).Returns(definition); + _materializerRegistry.IsMaterializerAvailable("materializer").Returns(true); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition, Arg.Any()).Returns(workflow); + _materializerRegistry.GetMaterializer("materializer").Returns(materializer); + _workflowGraphBuilder.BuildAsync(workflow, Arg.Any()).Returns(workflowGraph); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphAsync(filter); + + // Assert + Assert.NotNull(result); + Assert.Same(definition, result.WorkflowDefinition); + Assert.Same(workflowGraph, result.WorkflowGraph); + } + + [Fact] + public async Task TryFindWorkflowGraphsAsync_WithMultipleDefinitions_ReturnsAllResults() + { + // Arrange + var definition1 = TestHelpers.CreateWorkflowDefinition("def-1", "materializer"); + var definition2 = TestHelpers.CreateWorkflowDefinition("def-2", "unavailable-materializer"); + var definitions = new[] { definition1, definition2 }; + + var workflow1 = new Workflow { Identity = new("def-1", 1, "v1") }; + var graph1 = TestHelpers.CreateWorkflowGraph(workflow1); + + var filter = new WorkflowDefinitionFilter(); + _workflowDefinitionStore.FindManyAsync(filter, Arg.Any()).Returns(definitions); + + _materializerRegistry.IsMaterializerAvailable("materializer").Returns(true); + _materializerRegistry.IsMaterializerAvailable("unavailable-materializer").Returns(false); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition1, Arg.Any()).Returns(workflow1); + _materializerRegistry.GetMaterializer("materializer").Returns(materializer); + + _workflowGraphBuilder.BuildAsync(workflow1, Arg.Any()).Returns(graph1); + + var service = CreateService(); + + // Act + var result = await service.TryFindWorkflowGraphsAsync(filter); + + // Assert + var results = result.ToList(); + Assert.Equal(2, results.Count); + + var result1 = results.First(r => r.WorkflowDefinition?.DefinitionId == "def-1"); + Assert.Same(definition1, result1.WorkflowDefinition); + Assert.Same(graph1, result1.WorkflowGraph); + + var result2 = results.First(r => r.WorkflowDefinition?.DefinitionId == "def-2"); + Assert.Same(definition2, result2.WorkflowDefinition); + Assert.Null(result2.WorkflowGraph); + } + + private WorkflowDefinitionService CreateService() + { + return new( + _workflowDefinitionStore, + _workflowGraphBuilder, + _materializerRegistry, + _logger); + } + + /// + /// Sets up a materializer and graph builder for the given definition. + /// Returns the workflow and workflow graph that were configured. + /// + private (Workflow, WorkflowGraph) SetupMaterializerAndGraphBuilder(WorkflowDefinition definition) + { + var workflow = new Workflow { Identity = new(definition.DefinitionId, 1, "v1") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition, Arg.Any()).Returns(workflow); + + _materializerRegistry.GetMaterializer(definition.MaterializerName).Returns(materializer); + _workflowGraphBuilder.BuildAsync(workflow, Arg.Any()).Returns(workflowGraph); + + return (workflow, workflowGraph); + } + + /// + /// Sets up complete mocking chain for workflow graph retrieval including definition store, materializer, and graph builder. + /// Returns the workflow graph that was configured. + /// + private WorkflowGraph SetupCompleteWorkflowGraphMocking( + WorkflowDefinition definition, + bool materializerAvailable = true, + WorkflowDefinitionFilter? filter = null) + { + var workflow = new Workflow { Identity = new(definition.DefinitionId, 1, "v1") }; + var workflowGraph = TestHelpers.CreateWorkflowGraph(workflow); + + if (filter != null) + { + _workflowDefinitionStore.FindAsync(filter, Arg.Any()).Returns(definition); + } + else + { + _workflowDefinitionStore.FindAsync(Arg.Any(), Arg.Any()) + .Returns(definition); + } + + _materializerRegistry.IsMaterializerAvailable(definition.MaterializerName).Returns(materializerAvailable); + + if (materializerAvailable) + { + var materializer = Substitute.For(); + materializer.MaterializeAsync(definition, Arg.Any()).Returns(workflow); + _materializerRegistry.GetMaterializer(definition.MaterializerName).Returns(materializer); + _workflowGraphBuilder.BuildAsync(workflow, Arg.Any()).Returns(workflowGraph); + } + + return workflowGraph; + } +} diff --git a/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/BackgroundWorkflowCancellationDispatcherTests.cs b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/BackgroundWorkflowCancellationDispatcherTests.cs new file mode 100644 index 000000000..6e3947baf --- /dev/null +++ b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/BackgroundWorkflowCancellationDispatcherTests.cs @@ -0,0 +1,120 @@ +using Elsa.Common.Multitenancy; +using Elsa.Mediator; +using Elsa.Mediator.Contracts; +using Elsa.Mediator.Options; +using Elsa.Tenants.Mediator; +using Elsa.Workflows.Runtime.Commands; +using Elsa.Workflows.Runtime.Requests; +using NSubstitute; + +namespace Elsa.Workflows.Runtime.UnitTests.Services; + +public class BackgroundWorkflowCancellationDispatcherTests +{ + private readonly ICommandSender _commandSender = Substitute.For(); + private readonly ITenantAccessor _tenantAccessor = Substitute.For(); + + [Fact] + public async Task DispatchAsync_SendsCommandWithBackgroundStrategy() + { + // Arrange + var dispatcher = CreateDispatcher(); + var request = new DispatchCancelWorkflowRequest(); + + // Act + await dispatcher.DispatchAsync(request); + + // Assert + await _commandSender.Received(1).SendAsync( + Arg.Is(cmd => cmd.Request == request), + CommandStrategy.Background, + Arg.Any>(), + Arg.Any()); + } + + [Fact] + public async Task DispatchAsync_IncludesTenantIdInHeaders_WhenTenantIsPresent() + { + // Arrange + var tenantId = "test-tenant-123"; + var tenant = new Tenant + { + Id = tenantId + }; + _tenantAccessor.Tenant.Returns(tenant); + + var dispatcher = CreateDispatcher(); + var request = new DispatchCancelWorkflowRequest(); + + // Act + await dispatcher.DispatchAsync(request); + + // Assert + await _commandSender.Received(1).SendAsync( + Arg.Any(), + CommandStrategy.Background, + Arg.Is>(headers => + headers.ContainsKey(TenantHeaders.TenantIdKey) && + headers[TenantHeaders.TenantIdKey].ToString() == tenantId), + Arg.Any()); + } + + [Fact] + public async Task DispatchAsync_DoesNotIncludeTenantIdInHeaders_WhenTenantIsNull() + { + // Arrange + _tenantAccessor.Tenant.Returns((Tenant?)null); + + var dispatcher = CreateDispatcher(); + var request = new DispatchCancelWorkflowRequest(); + + // Act + await dispatcher.DispatchAsync(request); + + // Assert + await _commandSender.Received(1).SendAsync( + Arg.Any(), + CommandStrategy.Background, + Arg.Is>(headers => + !headers.ContainsKey(TenantHeaders.TenantIdKey)), + Arg.Any()); + } + + [Fact] + public async Task DispatchAsync_ReturnsResponse() + { + // Arrange + var dispatcher = CreateDispatcher(); + var request = new DispatchCancelWorkflowRequest(); + + // Act + var response = await dispatcher.DispatchAsync(request); + + // Assert + Assert.NotNull(response); + } + + [Fact] + public async Task DispatchAsync_PassesCancellationToken() + { + // Arrange + var dispatcher = CreateDispatcher(); + var request = new DispatchCancelWorkflowRequest(); + var cancellationToken = CancellationToken.None; + + // Act + await dispatcher.DispatchAsync(request, cancellationToken); + + // Assert + await _commandSender.Received(1).SendAsync( + Arg.Any(), + CommandStrategy.Background, + Arg.Any>(), + cancellationToken); + } + + private BackgroundWorkflowCancellationDispatcher CreateDispatcher() + { + return new(_commandSender, _tenantAccessor); + } +} \ No newline at end of file diff --git a/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/DefaultWorkflowDefinitionStorePopulatorTests.cs b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/DefaultWorkflowDefinitionStorePopulatorTests.cs index 5f960ea40..2bbcec480 100644 --- a/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/DefaultWorkflowDefinitionStorePopulatorTests.cs +++ b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/DefaultWorkflowDefinitionStorePopulatorTests.cs @@ -20,7 +20,7 @@ public class DefaultWorkflowDefinitionStorePopulatorTests _storeMock = Substitute.For(); _storeMock.FindManyAsync(Arg.Any(), Arg.Any()) .Returns(_workflowDefinitionsInStore); - _populator = new DefaultWorkflowDefinitionStorePopulator(() => new List(), + _populator = new(() => new List(), Substitute.For(), _storeMock, Substitute.For(), @@ -33,10 +33,10 @@ public class DefaultWorkflowDefinitionStorePopulatorTests [Fact(DisplayName = "When adding a new workflow it needs to be saved")] public async Task AddOrUpdateCoreAsync_NewWorkflowDefinition_AddsWorkflowDefinition() { - var workflow = new MaterializedWorkflow(new Workflow + var workflow = new MaterializedWorkflow(new() { - Identity = new WorkflowIdentity("a", 7, "1"), - Publication = new WorkflowPublication(true, true) + Identity = new("a", 7, "1"), + Publication = new(true, true) }, "Test", "Test"); await _populator.AddAsync(workflow); @@ -62,9 +62,9 @@ public class DefaultWorkflowDefinitionStorePopulatorTests } }); - var workflow = new MaterializedWorkflow(new Workflow + var workflow = new MaterializedWorkflow(new() { - Identity = new WorkflowIdentity("a", 1, "1"), + Identity = new("a", 1, "1"), Inputs = new List { new() @@ -94,10 +94,10 @@ public class DefaultWorkflowDefinitionStorePopulatorTests IsLatest = true, IsPublished = true }); - var workflow = new MaterializedWorkflow(new Workflow + var workflow = new MaterializedWorkflow(new() { - Identity = new WorkflowIdentity("a", 2, "2"), - Publication = new WorkflowPublication(workflowAddedIsLatest, workflowAddedIsPublished) + Identity = new("a", 2, "2"), + Publication = new(workflowAddedIsLatest, workflowAddedIsPublished) }, "Test", "Test"); await _populator.AddAsync(workflow); @@ -119,11 +119,11 @@ public class DefaultWorkflowDefinitionStorePopulatorTests IsPublished = true }); - var workflow = new MaterializedWorkflow(new Workflow + var workflow = new MaterializedWorkflow(new() { - Identity = new WorkflowIdentity("a", 3, "1"), + Identity = new("a", 3, "1"), Version = 1, - Publication = new WorkflowPublication(true, true) + Publication = new(true, true) }, "Test", "Test"); await _populator.AddAsync(workflow); @@ -143,9 +143,9 @@ public class DefaultWorkflowDefinitionStorePopulatorTests Version = 1, }); - var workflow = new MaterializedWorkflow(new Workflow + var workflow = new MaterializedWorkflow(new() { - Identity = new WorkflowIdentity("a", 2, "1") + Identity = new("a", 2, "1") }, "Test", "Test"); await _populator.AddAsync(workflow); @@ -166,10 +166,10 @@ public class DefaultWorkflowDefinitionStorePopulatorTests IsLatest = true }); - var workflow = new MaterializedWorkflow(new Workflow + var workflow = new MaterializedWorkflow(new() { - Identity = new WorkflowIdentity("a", 1, "1"), - Publication = new WorkflowPublication(true, true) + Identity = new("a", 1, "1"), + Publication = new(true, true) }, "Test", "Test"); await _populator.AddAsync(workflow); @@ -195,10 +195,10 @@ public class DefaultWorkflowDefinitionStorePopulatorTests } }); - var workflow = new MaterializedWorkflow(new Workflow + var workflow = new MaterializedWorkflow(new() { - Identity = new WorkflowIdentity("a", 1, "1"), - Publication = new WorkflowPublication(true, true) + Identity = new("a", 1, "1"), + Publication = new(true, true) }, "Test", "Test"); await _populator.AddAsync(workflow); diff --git a/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/LocalWorkflowClientTests.cs b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/LocalWorkflowClientTests.cs new file mode 100644 index 000000000..b795f8587 --- /dev/null +++ b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/LocalWorkflowClientTests.cs @@ -0,0 +1,131 @@ +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Exceptions; +using Elsa.Workflows.Management.Mappers; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Models; +using Elsa.Workflows.Runtime.Exceptions; +using Elsa.Workflows.Runtime.Messages; +using Microsoft.Extensions.Logging; +using NSubstitute; + +namespace Elsa.Workflows.Runtime.UnitTests.Services; + +public class LocalWorkflowClientTests +{ + private readonly IWorkflowInstanceManager _workflowInstanceManager = Substitute.For(); + private readonly IWorkflowDefinitionService _workflowDefinitionService = Substitute.For(); + private readonly IWorkflowRunner _workflowRunner = Substitute.For(); + private readonly IWorkflowCanceler _workflowCanceler = Substitute.For(); + private readonly WorkflowStateMapper _workflowStateMapper = Substitute.For(); + private readonly ILogger _logger = Substitute.For>(); + + [Fact] + public async Task CreateInstanceAsync_ThrowsWorkflowDefinitionNotFoundException_WhenDefinitionDoesNotExist() + { + // Arrange + var client = CreateClient(); + var definitionHandle = WorkflowDefinitionHandle.ByDefinitionId("non-existent-definition"); + var request = new CreateWorkflowInstanceRequest + { + WorkflowDefinitionHandle = definitionHandle + }; + + var findResult = new WorkflowGraphFindResult(null, null); + + _workflowDefinitionService.TryFindWorkflowGraphAsync(definitionHandle, Arg.Any()) + .Returns(findResult); + + // Act & Assert + await Assert.ThrowsAsync(() => + client.CreateInstanceAsync(request)); + } + + [Fact] + public async Task CreateInstanceAsync_ThrowsWorkflowMaterializerNotFoundException_WhenMaterializerNotAvailable() + { + // Arrange + var client = CreateClient(); + var definitionHandle = WorkflowDefinitionHandle.ByDefinitionId("test-definition"); + var request = new CreateWorkflowInstanceRequest + { + WorkflowDefinitionHandle = definitionHandle + }; + + var definition = new WorkflowDefinition + { + Id = "def-1", + DefinitionId = "test-definition", + MaterializerName = "unavailable-materializer" + }; + + var findResult = new WorkflowGraphFindResult(definition, null); // Null graph indicates materializer not available + + _workflowDefinitionService.TryFindWorkflowGraphAsync(definitionHandle, Arg.Any()) + .Returns(findResult); + + // Act & Assert + await Assert.ThrowsAsync(() => client.CreateInstanceAsync(request)); + } + + [Fact] + public async Task CreateAndRunInstanceAsync_ThrowsWorkflowDefinitionNotFoundException_WhenDefinitionDoesNotExist() + { + // Arrange + var client = CreateClient(); + var definitionHandle = WorkflowDefinitionHandle.ByDefinitionId("non-existent-definition"); + var request = new CreateAndRunWorkflowInstanceRequest + { + WorkflowDefinitionHandle = definitionHandle + }; + + var findResult = new WorkflowGraphFindResult(null, null); + + _workflowDefinitionService.TryFindWorkflowGraphAsync(definitionHandle, Arg.Any()) + .Returns(findResult); + + // Act & Assert + await Assert.ThrowsAsync(() => + client.CreateAndRunInstanceAsync(request)); + } + + [Fact] + public async Task CreateAndRunInstanceAsync_ThrowsWorkflowMaterializerNotFoundException_WhenMaterializerNotAvailable() + { + // Arrange + var client = CreateClient(); + var definitionHandle = WorkflowDefinitionHandle.ByDefinitionId("test-definition"); + var request = new CreateAndRunWorkflowInstanceRequest + { + WorkflowDefinitionHandle = definitionHandle + }; + + var definition = new WorkflowDefinition + { + Id = "def-1", + DefinitionId = "test-definition", + MaterializerName = "unavailable-materializer" + }; + + var findResult = new WorkflowGraphFindResult(definition, null); + + _workflowDefinitionService.TryFindWorkflowGraphAsync(definitionHandle, Arg.Any()) + .Returns(findResult); + + // Act & Assert + await Assert.ThrowsAsync(() => + client.CreateAndRunInstanceAsync(request)); + } + + private LocalWorkflowClient CreateClient() + { + return new( + "test-workflow-instance-id", + _workflowInstanceManager, + _workflowDefinitionService, + _workflowRunner, + _workflowCanceler, + _workflowStateMapper, + _logger); + } +}