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 <sipkeschoorstra@outlook.com>

* 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>
This commit is contained in:
Sipke Schoorstra 2026-01-19 08:59:12 +01:00 committed by GitHub
parent 40a7583f62
commit ca88051573
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
32 changed files with 1986 additions and 118 deletions

View file

@ -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; }
}

View file

@ -25,8 +25,13 @@ public interface ICacheManager
/// </summary>
ValueTask TriggerTokenAsync(string key, CancellationToken cancellationToken = default);
/// <summary>
/// Finds an item from the cache, or creates it if it doesn't exist.
/// </summary>
Task<TItem?> FindOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory);
/// <summary>
/// Gets an item from the cache, or creates it if it doesn't exist.
/// </summary>
Task<TItem?> GetOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory);
Task<TItem> GetOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory);
}

View file

@ -24,8 +24,21 @@ public class CacheManager(IMemoryCache memoryCache, IChangeTokenSignaler changeT
}
/// <inheritdoc />
public async Task<TItem?> GetOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory)
public async Task<TItem?> FindOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory)
{
return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry));
}
/// <summary>
/// Retrieves a cached item by the specified key or creates a new one using the provided factory function.
/// </summary>
/// <param name="key">The key used to identify the cached item.</param>
/// <param name="factory">A factory function that provides the value to be cached if it does not already exist.</param>
/// <typeparam name="TItem">The type of the item to retrieve or create.</typeparam>
/// <returns>The cached or newly created item.</returns>
/// <exception cref="InvalidOperationException">Thrown if the factory function returns null.</exception>
public async Task<TItem> GetOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory)
{
return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry)) ?? throw new InvalidOperationException($"Factory returned null for cache key: {key}.");
}
}

View file

@ -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);

View file

@ -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);

View file

@ -1,4 +1,3 @@
using Elsa.Dsl.ElsaScript.Compiler;
using Elsa.Dsl.ElsaScript.Features;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;

View file

@ -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<Request>
internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefinitionLinker linker, IMaterializerRegistry materializerRegistry) : ElsaEndpoint<Request>
{
public override void Configure()
{
@ -27,6 +28,13 @@ internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefini
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);
}

View file

@ -6,4 +6,9 @@ namespace Elsa.Workflows.Api.Models;
public class LinkedWorkflowDefinitionSummary : WorkflowDefinitionSummary
{
public Link[]? Links { get; set; }
/// <summary>
/// Indicates whether the materializer for this workflow definition is currently available.
/// </summary>
public bool IsMaterializerAvailable { get; set; }
}

View file

@ -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;
/// <inheritdoc/>
public class StaticWorkflowDefinitionLinker(
IOptions<ManagementOptions> managementOptions,
WorkflowDefinitionMapper workflowDefinitionMapper)
WorkflowDefinitionMapper workflowDefinitionMapper,
IMaterializerRegistry materializerRegistry)
: IWorkflowDefinitionLinker
{
/// <inheritdoc />
@ -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<LinkedWorkflowDefinitionSummary>
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();

View file

@ -0,0 +1,26 @@
namespace Elsa.Workflows.Management;
/// <summary>
/// Provides access to registered workflow materializers and their availability.
/// </summary>
public interface IMaterializerRegistry
{
/// <summary>
/// Gets all registered materializers.
/// </summary>
IEnumerable<IWorkflowMaterializer> GetMaterializers();
/// <summary>
/// Gets a materializer by name.
/// </summary>
/// <param name="name">The name of the materializer.</param>
/// <returns>The materializer, or null if not found.</returns>
IWorkflowMaterializer? GetMaterializer(string name);
/// <summary>
/// Checks if a materializer with the specified name is available.
/// </summary>
/// <param name="name">The name of the materializer.</param>
/// <returns>True if the materializer is available, false otherwise.</returns>
bool IsMaterializerAvailable(string name);
}

View file

@ -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 <see cref="WorkflowGraph"/>s that match the specified <see cref="WorkflowDefinitionFilter"/>.
/// </summary>
Task<IEnumerable<WorkflowGraph>> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
/// <summary>
/// Attempts to retrieve a <see cref="WorkflowGraph"/> and its corresponding <see cref="WorkflowDefinition"/>
/// based on the specified workflow definition ID and version options.
/// </summary>
/// <param name="definitionId">The ID of the workflow definition to search for.</param>
/// <param name="versionOptions">The versioning options that determine which version of the workflow definition to search for.</param>
/// <param name="cancellationToken">A token to monitor for cancellation requests.</param>
/// <returns>
/// A <see cref="WorkflowGraphFindResult"/> containing the workflow graph and its associated definition,
/// or indicating that they do not exist.
/// </returns>
Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default);
/// <summary>
/// Attempts to locate a <see cref="WorkflowGraph"/> based on the specified workflow definition version ID.
/// </summary>
/// <param name="definitionVersionId">The unique identifier of the workflow definition version.</param>
/// <param name="cancellationToken">A token to monitor for cancellation requests.</param>
/// <returns>
/// A <see cref="WorkflowGraphFindResult"/> containing the workflow graph and its associated definition,
/// or indicating that they do not exist.
/// </returns>
Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default);
/// <summary>
/// Attempts to retrieve the workflow graph corresponding to the specified <see cref="WorkflowDefinitionHandle"/>.
/// </summary>
/// <param name="definitionHandle">The handle identifying the workflow definition and its version.</param>
/// <param name="cancellationToken">A token to cancel the asynchronous operation.</param>
/// <returns>
/// A <see cref="WorkflowGraphFindResult"/> containing the workflow graph and its associated definition,
/// or indicating that they do not exist.
/// </returns>
Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default);
/// <summary>
/// Attempts to retrieve a <see cref="WorkflowGraph"/> and its associated <see cref="WorkflowDefinition"/>
/// for the specified workflow definition using the provided filter criteria.
/// </summary>
/// <param name="filter">The criteria used to filter and locate the workflow definition and graph.</param>
/// <param name="cancellationToken">A token to monitor for cancellation requests.</param>
/// <returns>
/// A <see cref="WorkflowGraphFindResult"/> containing the workflow graph and its associated definition,
/// or indicating that they do not exist.
/// </returns>
Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
/// <summary>
/// Attempts to find and retrieve a collection of <see cref="WorkflowGraphFindResult"/> objects based on the specified <see cref="WorkflowDefinitionFilter"/>.
/// </summary>
/// <param name="filter">The filter specifying criteria to identify the workflow definitions to search for.</param>
/// <param name="cancellationToken">A token to observe while waiting for the task to complete.</param>
/// <returns>
/// A collection of <see cref="WorkflowGraphFindResult"/> 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.
/// </returns>
Task<IEnumerable<WorkflowGraphFindResult>> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
}

View file

@ -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;
}

View file

@ -255,6 +255,7 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module)
.AddScoped<WorkflowDefinitionActivityDescriptorFactory>()
.AddScoped<WorkflowDefinitionActivityProvider>()
.AddScoped<IWorkflowDefinitionActivityRegistryUpdater, WorkflowDefinitionActivityRegistryUpdater>()
.AddScoped<IMaterializerRegistry, MaterializerRegistry>()
.AddScoped<IWorkflowDefinitionService, WorkflowDefinitionService>()
.AddScoped<IWorkflowSerializer, WorkflowSerializer>()
.AddScoped<IWorkflowValidator, WorkflowValidator>()

View file

@ -12,7 +12,7 @@ public class TimestampFilter
/// <summary>
/// Gets or sets the column to filter by.
/// </summary>
public required string Column { get; set; } = default!;
public required string Column { get; set; } = null!;
/// <summary>
/// Gets or sets the operator to use for filtering.

View file

@ -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;
}

View file

@ -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 <see cref="IWorkflowDefinitionService"/> with caching capabilities.
/// </summary>
[UsedImplicitly]
public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowDefinitionService
public class CachingWorkflowDefinitionService(
IWorkflowDefinitionService decoratedService,
IWorkflowDefinitionCacheManager cacheManager,
IWorkflowDefinitionStore workflowDefinitionStore,
IMaterializerRegistry materializerRegistry) : IWorkflowDefinitionService
{
/// <inheritdoc />
public async Task<WorkflowGraph> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
public Task<WorkflowGraph> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
return await decoratedService.MaterializeWorkflowAsync(definition, cancellationToken);
return decoratedService.MaterializeWorkflowAsync(definition, cancellationToken);
}
/// <inheritdoc />
@ -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<WorkflowDefinition?> 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
/// <inheritdoc />
public Task<WorkflowDefinition?> 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<WorkflowDefinition?> 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<WorkflowGraph?> 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<WorkflowGraph?> 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
/// <inheritdoc />
public Task<WorkflowGraph?> 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<WorkflowGraph?> 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);
}
/// <inheritdoc />
public async Task<IEnumerable<WorkflowGraph>> 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<T?> GetFromCacheAsync<T>(string cacheKey, Func<Task<T?>> getObjectFunc, Func<T, string> getChangeTokenKeyFunc)
/// <inheritdoc />
public async Task<WorkflowGraphFindResult> 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;
}
/// <inheritdoc />
public async Task<WorkflowGraphFindResult> 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;
}
/// <inheritdoc />
public Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default)
{
var filter = definitionHandle.ToFilter();
return TryFindWorkflowGraphAsync(filter, cancellationToken);
}
/// <inheritdoc />
public async Task<WorkflowGraphFindResult> 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;
}
/// <inheritdoc />
public async Task<IEnumerable<WorkflowGraphFindResult>> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken);
var results = new List<WorkflowGraphFindResult>();
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<T?> FindFromCacheAsync<T>(string cacheKey, Func<Task<T?>> getObjectFunc, Func<T, string> getChangeTokenKeyFunc) where T : class
{
return await GetFromCacheAsync(
cacheKey,
getObjectFunc,
obj => obj != null ? getChangeTokenKeyFunc(obj) : null);
}
private async Task<T> GetFromCacheAsync<T>(string cacheKey, Func<Task<T>> getObjectFunc, Func<T, string?> 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;
});
}

View file

@ -0,0 +1,22 @@
namespace Elsa.Workflows.Management.Services;
/// <inheritdoc />
public class MaterializerRegistry(Func<IEnumerable<IWorkflowMaterializer>> materializers) : IMaterializerRegistry
{
private readonly Lazy<IReadOnlyCollection<IWorkflowMaterializer>> _materializers = new(() => materializers().ToArray());
/// <inheritdoc />
public IEnumerable<IWorkflowMaterializer> GetMaterializers() => _materializers.Value;
/// <inheritdoc />
public IWorkflowMaterializer? GetMaterializer(string name)
{
return _materializers.Value.FirstOrDefault(x => x.Name == name);
}
/// <inheritdoc />
public bool IsMaterializerAvailable(string name)
{
return _materializers.Value.Any(x => x.Name == name);
}
}

View file

@ -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<IEnumerable<IWorkflowMaterializer>> materializers)
IMaterializerRegistry materializerRegistry,
ILogger<WorkflowDefinitionService> logger)
: IWorkflowDefinitionService
{
/// <inheritdoc />
public async Task<WorkflowGraph> 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<WorkflowGraph?> 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);
}
/// <inheritdoc />
public async Task<WorkflowGraph?> 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);
}
/// <inheritdoc />
public async Task<WorkflowGraph?> 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);
}
/// <inheritdoc />
public async Task<WorkflowGraph?> 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);
}
/// <inheritdoc />
@ -109,4 +96,81 @@ public class WorkflowDefinitionService(
return workflowGraphs;
}
public async Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken);
return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
}
public async Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken);
return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
}
public async Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionHandle, cancellationToken);
return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
}
public async Task<WorkflowGraphFindResult> TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken);
return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
}
public async Task<IEnumerable<WorkflowGraphFindResult>> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken);
var results = new List<WorkflowGraphFindResult>();
foreach (var workflowDefinition in workflowDefinitions)
{
var result = await MaterializeWorkflowGraphFindResultAsync(workflowDefinition, cancellationToken);
results.Add(result);
}
return results;
}
/// <summary>
/// Attempts to materialize a workflow graph from the given workflow definition if a suitable materializer is available.
/// </summary>
/// <param name="definition">The workflow definition to materialize. Can be null.</param>
/// <param name="cancellationToken">A token to monitor for cancellation requests.</param>
/// <returns>
/// A <see cref="WorkflowGraph"/> if materialization is successful; otherwise, null.
/// </returns>
private async Task<WorkflowGraph?> 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;
}
/// <summary>
/// Attempts to materialize a workflow graph from the given workflow definition.
/// </summary>
/// <param name="definition">The workflow definition to materialize the graph for. May be null.</param>
/// <param name="cancellationToken">A token to observe while waiting for the task to complete.</param>
/// <returns>A result containing the materialized workflow graph and its corresponding workflow definition, if successfully materialized; otherwise, returns the definition with a null graph.</returns>
private async Task<WorkflowGraphFindResult> 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);
}
}

View file

@ -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);

View file

@ -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;
}

View file

@ -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;

View file

@ -1,5 +1,7 @@
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;
@ -9,13 +11,18 @@ namespace Elsa.Workflows.Runtime;
/// <summary>
/// Dispatches workflow cancellation requests to a local background worker.
/// </summary>
public class BackgroundWorkflowCancellationDispatcher(ICommandSender commandSender) : IWorkflowCancellationDispatcher
public class BackgroundWorkflowCancellationDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IWorkflowCancellationDispatcher
{
/// <inheritdoc />
public async Task<DispatchCancelWorkflowsResponse> DispatchAsync(DispatchCancelWorkflowRequest request, CancellationToken cancellationToken = default)
{
var command = new CancelWorkflowsCommand(request);
await commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken);
return new DispatchCancelWorkflowsResponse();
await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken);
return new();
}
private IDictionary<object, object> CreateHeaders()
{
return TenantHeaders.CreateHeaders(tenantAccessor.Tenant?.Id);
}
}

View file

@ -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<WorkflowGraph> 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!;
}
}

View file

@ -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);

View file

@ -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);
}
}

View file

@ -0,0 +1,40 @@
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.UnitTests.Helpers;
/// <summary>
/// Provides helper methods for creating test data in Workflow Management unit tests.
/// </summary>
public static class TestHelpers
{
/// <summary>
/// Creates a workflow definition with the specified parameters.
/// </summary>
public static WorkflowDefinition CreateWorkflowDefinition(string definitionId, string materializerName)
{
return new WorkflowDefinition
{
DefinitionId = definitionId,
MaterializerName = materializerName,
Version = 1
};
}
/// <summary>
/// Creates a workflow graph from a workflow, ensuring proper initialization.
/// </summary>
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 });
}
}

View file

@ -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<IWorkflowDefinitionService>();
private readonly IWorkflowDefinitionCacheManager _cacheManager = Substitute.For<IWorkflowDefinitionCacheManager>();
private readonly IWorkflowDefinitionStore _workflowDefinitionStore = Substitute.For<IWorkflowDefinitionStore>();
private readonly IMaterializerRegistry _materializerRegistry = Substitute.For<IMaterializerRegistry>();
private readonly ICacheManager _cache = Substitute.For<ICacheManager>();
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<CancellationToken>()).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<CancellationToken>());
}
[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<CancellationToken>())
.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<CancellationToken>()).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<WorkflowDefinitionFilter>()).Returns(cacheKey);
_decoratedService.FindWorkflowDefinitionAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(definition);
var service = CreateService();
// Act
var result = await service.FindWorkflowDefinitionAsync(handle);
// Assert
Assert.Same(definition, result);
await _decoratedService.Received(1).FindWorkflowDefinitionAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>());
}
[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<CancellationToken>()).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<CancellationToken>())
.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<CancellationToken>()).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<WorkflowDefinitionFilter>()).Returns(cacheKey);
_decoratedService.FindWorkflowGraphAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(workflowGraph);
var service = CreateService();
// Act
var result = await service.FindWorkflowGraphAsync(handle);
// Assert
Assert.Same(workflowGraph, result);
await _decoratedService.Received(1).FindWorkflowGraphAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>());
}
[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<CancellationToken>()).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<CancellationToken>()).Returns(definitions);
_cacheManager.CreateWorkflowVersionCacheKey("id-1").Returns("cache-key-1");
_cacheManager.CreateWorkflowVersionCacheKey("id-2").Returns("cache-key-2");
_decoratedService.MaterializeWorkflowAsync(definition1, Arg.Any<CancellationToken>()).Returns(graph1);
_decoratedService.MaterializeWorkflowAsync(definition2, Arg.Any<CancellationToken>()).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<CancellationToken>())
.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<CancellationToken>()).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<WorkflowDefinitionFilter>()).Returns(cacheKey);
_decoratedService.TryFindWorkflowGraphAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(findResult);
var service = CreateService();
// Act
var result = await service.TryFindWorkflowGraphAsync(handle);
// Assert
Assert.Same(findResult, result);
await _decoratedService.Received(1).TryFindWorkflowGraphAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>());
}
[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<CancellationToken>()).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<CancellationToken>()).Returns(definitions);
_cacheManager.CreateWorkflowVersionCacheKey("id-1").Returns("cache-key-1");
_materializerRegistry.IsMaterializerAvailable("materializer").Returns(true);
_decoratedService.MaterializeWorkflowAsync(definition, Arg.Any<CancellationToken>()).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<CancellationToken>()).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<WorkflowDefinition>(), Arg.Any<CancellationToken>());
}
[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<CancellationToken>()).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<CancellationToken>()).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);
}
/// <summary>
/// Sets up the cache manager with proper mock behavior for all cache operations.
/// </summary>
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<WorkflowDefinition>(Arg.Any<object>(), Arg.Any<Func<ICacheEntry, Task<WorkflowDefinition>>>())
.Returns(async callInfo => await callInfo.Arg<Func<ICacheEntry, Task<WorkflowDefinition>>>()(Substitute.For<ICacheEntry>()));
_cache.GetOrCreateAsync<WorkflowGraph>(Arg.Any<object>(), Arg.Any<Func<ICacheEntry, Task<WorkflowGraph>>>())
.Returns(async callInfo => await callInfo.Arg<Func<ICacheEntry, Task<WorkflowGraph>>>()(Substitute.For<ICacheEntry>()));
_cache.GetOrCreateAsync<WorkflowGraphFindResult>(Arg.Any<object>(), Arg.Any<Func<ICacheEntry, Task<WorkflowGraphFindResult>>>())
.Returns(async callInfo => await callInfo.Arg<Func<ICacheEntry, Task<WorkflowGraphFindResult>>>()(Substitute.For<ICacheEntry>()));
_cache.FindOrCreateAsync<WorkflowDefinition>(Arg.Any<object>(), Arg.Any<Func<ICacheEntry, Task<WorkflowDefinition>>>())
.Returns(async callInfo => await callInfo.Arg<Func<ICacheEntry, Task<WorkflowDefinition>>>()(Substitute.For<ICacheEntry>()));
_cache.FindOrCreateAsync<WorkflowGraph>(Arg.Any<object>(), Arg.Any<Func<ICacheEntry, Task<WorkflowGraph>>>())
.Returns(async callInfo => await callInfo.Arg<Func<ICacheEntry, Task<WorkflowGraph>>>()(Substitute.For<ICacheEntry>()));
}
/// <summary>
/// Creates a workflow and workflow graph for testing.
/// </summary>
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);
}
/// <summary>
/// Creates a workflow graph find result for testing.
/// </summary>
private WorkflowGraphFindResult CreateWorkflowGraphFindResult(string definitionId = "def-1", string materializerName = "materializer")
{
var definition = TestHelpers.CreateWorkflowDefinition(definitionId, materializerName);
var (_, workflowGraph) = CreateWorkflowAndGraph(definitionId);
return new(definition, workflowGraph);
}
}

View file

@ -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<IWorkflowMaterializer> 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<IWorkflowMaterializer>();
materializer.Name.Returns(name);
return materializer;
}
private static MaterializerRegistry CreateRegistry(params IWorkflowMaterializer[] materializers)
{
return new(() => materializers);
}
}

View file

@ -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<IWorkflowDefinitionStore>();
private readonly IWorkflowGraphBuilder _workflowGraphBuilder = Substitute.For<IWorkflowGraphBuilder>();
private readonly IMaterializerRegistry _materializerRegistry = Substitute.For<IMaterializerRegistry>();
private readonly ILogger<WorkflowDefinitionService> _logger = Substitute.For<ILogger<WorkflowDefinitionService>>();
[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<WorkflowMaterializerNotFoundException>(() => service.MaterializeWorkflowAsync(definition));
}
[Fact]
public async Task FindWorkflowDefinitionAsync_ByDefinitionIdAndVersionOptions_ReturnsDefinition()
{
// Arrange
var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer");
_workflowDefinitionStore.FindAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.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<WorkflowDefinitionFilter>(),
Arg.Any<CancellationToken>());
}
[Fact]
public async Task FindWorkflowDefinitionAsync_ByDefinitionVersionId_ReturnsDefinition()
{
// Arrange
var definition = TestHelpers.CreateWorkflowDefinition("def-1", "materializer");
definition.Id = "version-id-1";
_workflowDefinitionStore.FindAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.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<WorkflowDefinitionFilter>(f => f.Id == "version-id-1"),
Arg.Any<CancellationToken>());
}
[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<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(definition);
var service = CreateService();
// Act
var result = await service.FindWorkflowDefinitionAsync(handle);
// Assert
Assert.Same(definition, result);
await _workflowDefinitionStore.Received(1).FindAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>());
}
[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<CancellationToken>()).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<CancellationToken>());
}
[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<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.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<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(definition);
_materializerRegistry.IsMaterializerAvailable("materializer").Returns(true);
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition, Arg.Any<CancellationToken>()).Returns(workflow);
_materializerRegistry.GetMaterializer("materializer").Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow, Arg.Any<CancellationToken>()).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<CancellationToken>()).Returns(definition);
_materializerRegistry.IsMaterializerAvailable("materializer").Returns(true);
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition, Arg.Any<CancellationToken>()).Returns(workflow);
_materializerRegistry.GetMaterializer("materializer").Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow, Arg.Any<CancellationToken>()).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<CancellationToken>()).Returns(definitions);
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition1, Arg.Any<CancellationToken>()).Returns(workflow1);
materializer.MaterializeAsync(definition2, Arg.Any<CancellationToken>()).Returns(workflow2);
_materializerRegistry.GetMaterializer("materializer").Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow1, Arg.Any<CancellationToken>()).Returns(graph1);
_workflowGraphBuilder.BuildAsync(workflow2, Arg.Any<CancellationToken>()).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<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(definition);
_materializerRegistry.IsMaterializerAvailable("materializer").Returns(true);
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition, Arg.Any<CancellationToken>()).Returns(workflow);
_materializerRegistry.GetMaterializer("materializer").Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow, Arg.Any<CancellationToken>()).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<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.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<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(definition);
_materializerRegistry.IsMaterializerAvailable("materializer").Returns(true);
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition, Arg.Any<CancellationToken>()).Returns(workflow);
_materializerRegistry.GetMaterializer("materializer").Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow, Arg.Any<CancellationToken>()).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<CancellationToken>()).Returns(definition);
_materializerRegistry.IsMaterializerAvailable("materializer").Returns(true);
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition, Arg.Any<CancellationToken>()).Returns(workflow);
_materializerRegistry.GetMaterializer("materializer").Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow, Arg.Any<CancellationToken>()).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<CancellationToken>()).Returns(definitions);
_materializerRegistry.IsMaterializerAvailable("materializer").Returns(true);
_materializerRegistry.IsMaterializerAvailable("unavailable-materializer").Returns(false);
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition1, Arg.Any<CancellationToken>()).Returns(workflow1);
_materializerRegistry.GetMaterializer("materializer").Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow1, Arg.Any<CancellationToken>()).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);
}
/// <summary>
/// Sets up a materializer and graph builder for the given definition.
/// Returns the workflow and workflow graph that were configured.
/// </summary>
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<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition, Arg.Any<CancellationToken>()).Returns(workflow);
_materializerRegistry.GetMaterializer(definition.MaterializerName).Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow, Arg.Any<CancellationToken>()).Returns(workflowGraph);
return (workflow, workflowGraph);
}
/// <summary>
/// Sets up complete mocking chain for workflow graph retrieval including definition store, materializer, and graph builder.
/// Returns the workflow graph that was configured.
/// </summary>
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<CancellationToken>()).Returns(definition);
}
else
{
_workflowDefinitionStore.FindAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(definition);
}
_materializerRegistry.IsMaterializerAvailable(definition.MaterializerName).Returns(materializerAvailable);
if (materializerAvailable)
{
var materializer = Substitute.For<IWorkflowMaterializer>();
materializer.MaterializeAsync(definition, Arg.Any<CancellationToken>()).Returns(workflow);
_materializerRegistry.GetMaterializer(definition.MaterializerName).Returns(materializer);
_workflowGraphBuilder.BuildAsync(workflow, Arg.Any<CancellationToken>()).Returns(workflowGraph);
}
return workflowGraph;
}
}

View file

@ -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<ICommandSender>();
private readonly ITenantAccessor _tenantAccessor = Substitute.For<ITenantAccessor>();
[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<CancelWorkflowsCommand>(cmd => cmd.Request == request),
CommandStrategy.Background,
Arg.Any<IDictionary<object, object>>(),
Arg.Any<CancellationToken>());
}
[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<CancelWorkflowsCommand>(),
CommandStrategy.Background,
Arg.Is<IDictionary<object, object>>(headers =>
headers.ContainsKey(TenantHeaders.TenantIdKey) &&
headers[TenantHeaders.TenantIdKey].ToString() == tenantId),
Arg.Any<CancellationToken>());
}
[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<CancelWorkflowsCommand>(),
CommandStrategy.Background,
Arg.Is<IDictionary<object, object>>(headers =>
!headers.ContainsKey(TenantHeaders.TenantIdKey)),
Arg.Any<CancellationToken>());
}
[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<CancelWorkflowsCommand>(),
CommandStrategy.Background,
Arg.Any<IDictionary<object, object>>(),
cancellationToken);
}
private BackgroundWorkflowCancellationDispatcher CreateDispatcher()
{
return new(_commandSender, _tenantAccessor);
}
}

View file

@ -20,7 +20,7 @@ public class DefaultWorkflowDefinitionStorePopulatorTests
_storeMock = Substitute.For<IWorkflowDefinitionStore>();
_storeMock.FindManyAsync(Arg.Any<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(_workflowDefinitionsInStore);
_populator = new DefaultWorkflowDefinitionStorePopulator(() => new List<IWorkflowsProvider>(),
_populator = new(() => new List<IWorkflowsProvider>(),
Substitute.For<ITriggerIndexer>(),
_storeMock,
Substitute.For<IActivitySerializer>(),
@ -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<InputDefinition>
{
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);

View file

@ -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<IWorkflowInstanceManager>();
private readonly IWorkflowDefinitionService _workflowDefinitionService = Substitute.For<IWorkflowDefinitionService>();
private readonly IWorkflowRunner _workflowRunner = Substitute.For<IWorkflowRunner>();
private readonly IWorkflowCanceler _workflowCanceler = Substitute.For<IWorkflowCanceler>();
private readonly WorkflowStateMapper _workflowStateMapper = Substitute.For<WorkflowStateMapper>();
private readonly ILogger<LocalWorkflowClient> _logger = Substitute.For<ILogger<LocalWorkflowClient>>();
[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<CancellationToken>())
.Returns(findResult);
// Act & Assert
await Assert.ThrowsAsync<WorkflowDefinitionNotFoundException>(() =>
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<CancellationToken>())
.Returns(findResult);
// Act & Assert
await Assert.ThrowsAsync<WorkflowMaterializerNotFoundException>(() => 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<CancellationToken>())
.Returns(findResult);
// Act & Assert
await Assert.ThrowsAsync<WorkflowDefinitionNotFoundException>(() =>
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<CancellationToken>())
.Returns(findResult);
// Act & Assert
await Assert.ThrowsAsync<WorkflowMaterializerNotFoundException>(() =>
client.CreateAndRunInstanceAsync(request));
}
private LocalWorkflowClient CreateClient()
{
return new(
"test-workflow-instance-id",
_workflowInstanceManager,
_workflowDefinitionService,
_workflowRunner,
_workflowCanceler,
_workflowStateMapper,
_logger);
}
}