diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs
index fe1a4cc55..dcbc2ed55 100644
--- a/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs
+++ b/src/clients/Elsa.Api.Client/Resources/WorkflowDefinitions/Models/WorkflowDefinitionSummary.cs
@@ -9,8 +9,9 @@ namespace Elsa.Api.Client.Resources.WorkflowDefinitions.Models;
[PublicAPI]
public class WorkflowDefinitionSummary : LinkedEntity
{
- public string DefinitionId { get; set; } = default!;
- public string Name { get; set; } = default!;
+ public string DefinitionId { get; set; } = null!;
+ public string Name { get; set; } = null!;
public string? Description { get; set; }
- public string MaterializerName { get; set; } = default!;
+ public string MaterializerName { get; set; } = null!;
+ public bool IsMaterializerAvailable { get; set; }
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Caching/Contracts/ICacheManager.cs b/src/modules/Elsa.Caching/Contracts/ICacheManager.cs
index 6e93fe8c8..1b75ffbed 100644
--- a/src/modules/Elsa.Caching/Contracts/ICacheManager.cs
+++ b/src/modules/Elsa.Caching/Contracts/ICacheManager.cs
@@ -25,8 +25,13 @@ public interface ICacheManager
///
ValueTask TriggerTokenAsync(string key, CancellationToken cancellationToken = default);
+ ///
+ /// Finds an item from the cache, or creates it if it doesn't exist.
+ ///
+ Task FindOrCreateAsync(object key, Func> factory);
+
///
/// Gets an item from the cache, or creates it if it doesn't exist.
///
- Task GetOrCreateAsync(object key, Func> factory);
+ Task GetOrCreateAsync(object key, Func> factory);
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Caching/Services/CacheManager.cs b/src/modules/Elsa.Caching/Services/CacheManager.cs
index 6db501327..657a1c66d 100644
--- a/src/modules/Elsa.Caching/Services/CacheManager.cs
+++ b/src/modules/Elsa.Caching/Services/CacheManager.cs
@@ -24,8 +24,21 @@ public class CacheManager(IMemoryCache memoryCache, IChangeTokenSignaler changeT
}
///
- public async Task GetOrCreateAsync(object key, Func> factory)
+ public async Task FindOrCreateAsync(object key, Func> factory)
{
return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry));
}
+
+ ///
+ /// Retrieves a cached item by the specified key or creates a new one using the provided factory function.
+ ///
+ /// The key used to identify the cached item.
+ /// A factory function that provides the value to be cached if it does not already exist.
+ /// The type of the item to retrieve or create.
+ /// The cached or newly created item.
+ /// Thrown if the factory function returns null.
+ public async Task GetOrCreateAsync(object key, Func> factory)
+ {
+ return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry)) ?? throw new InvalidOperationException($"Factory returned null for cache key: {key}.");
+ }
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Common/Codecs/Zstd.cs b/src/modules/Elsa.Common/Codecs/Zstd.cs
index 11a49c5d6..c8b46785a 100644
--- a/src/modules/Elsa.Common/Codecs/Zstd.cs
+++ b/src/modules/Elsa.Common/Codecs/Zstd.cs
@@ -15,7 +15,7 @@ public class Zstd : ICompressionCodec
{
var inputBytes = Encoding.UTF8.GetBytes(input);
var span = inputBytes.AsSpan();
- var result = Iron.Compress(Codec.Zstd, span);
+ using var result = Iron.Compress(Codec.Zstd, span);
var compressedBytes = result.AsSpan();
var compressedString = Convert.ToBase64String(compressedBytes);
@@ -27,7 +27,7 @@ public class Zstd : ICompressionCodec
{
var inputBytes = Convert.FromBase64String(input);
var span = inputBytes.AsSpan();
- var result = Iron.Decompress(Codec.Zstd, span);
+ using var result = Iron.Decompress(Codec.Zstd, span);
var decompressedBytes = result.AsSpan();
var decompressedString = Encoding.UTF8.GetString(decompressedBytes);
diff --git a/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs
index b4d2e63ac..1667c45c5 100644
--- a/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs
+++ b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs
@@ -21,7 +21,7 @@ public class CachingHttpWorkflowLookupService(
var tenantIdPrefix = !string.IsNullOrEmpty(tenantId) ? $"{tenantId}:" : string.Empty;
var key = $"{tenantIdPrefix}http-workflow:{bookmarkHash}";
var cache = cacheManager.Cache;
- return await cache.GetOrCreateAsync(key, async entry =>
+ return await cache.FindOrCreateAsync(key, async entry =>
{
var cachingOptions = cache.CachingOptions.Value;
entry.SetSlidingExpiration(cachingOptions.CacheDuration);
diff --git a/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs b/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs
index 5c18cf948..67236da2c 100644
--- a/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs
+++ b/src/modules/Elsa.WorkflowProviders.BlobStorage.ElsaScript/Features/ElsaScriptBlobStorageFeature.cs
@@ -1,4 +1,3 @@
-using Elsa.Dsl.ElsaScript.Compiler;
using Elsa.Dsl.ElsaScript.Features;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs
index 813ca17e7..93f333056 100644
--- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs
+++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/GetByDefinitionId/Endpoint.cs
@@ -3,11 +3,12 @@ using Elsa.Common.Models;
using Elsa.Workflows.Management;
using Elsa.Workflows.Models;
using JetBrains.Annotations;
+using Microsoft.AspNetCore.Http;
namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.GetByDefinitionId;
[PublicAPI]
-internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefinitionLinker linker) : ElsaEndpoint
+internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefinitionLinker linker, IMaterializerRegistry materializerRegistry) : ElsaEndpoint
{
public override void Configure()
{
@@ -20,13 +21,20 @@ internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefini
var versionOptions = request.VersionOptions != null ? VersionOptions.FromString(request.VersionOptions) : VersionOptions.Latest;
var filter = WorkflowDefinitionHandle.ByDefinitionId(request.DefinitionId, versionOptions).ToFilter();
var definition = await store.FindAsync(filter, cancellationToken);
-
+
if (definition == null)
{
await Send.NotFoundAsync(cancellationToken);
return;
}
+ if (!materializerRegistry.IsMaterializerAvailable(definition.MaterializerName))
+ {
+ AddError($"The workflow materializer '{definition.MaterializerName}' is not available. The materializer may be disabled or not registered.");
+ await Send.ErrorsAsync(StatusCodes.Status422UnprocessableEntity, cancellationToken);
+ return;
+ }
+
var model = await linker.MapAsync(definition, cancellationToken);
await Send.OkAsync(model, cancellationToken);
}
diff --git a/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs b/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs
index ebf28f35f..79a56cd68 100644
--- a/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs
+++ b/src/modules/Elsa.Workflows.Api/Models/LinkedWorkflowDefinitionSummary.cs
@@ -6,4 +6,9 @@ namespace Elsa.Workflows.Api.Models;
public class LinkedWorkflowDefinitionSummary : WorkflowDefinitionSummary
{
public Link[]? Links { get; set; }
+
+ ///
+ /// Indicates whether the materializer for this workflow definition is currently available.
+ ///
+ public bool IsMaterializerAvailable { get; set; }
}
diff --git a/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs b/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs
index da2a3351e..87ebc22ff 100644
--- a/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs
+++ b/src/modules/Elsa.Workflows.Api/Services/StaticWorkflowDefinitionLinker.cs
@@ -1,5 +1,6 @@
using Elsa.Models;
using Elsa.Workflows.Api.Models;
+using Elsa.Workflows.Management;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Mappers;
using Elsa.Workflows.Management.Models;
@@ -11,7 +12,8 @@ namespace Elsa.Workflows.Api.Services;
///
public class StaticWorkflowDefinitionLinker(
IOptions managementOptions,
- WorkflowDefinitionMapper workflowDefinitionMapper)
+ WorkflowDefinitionMapper workflowDefinitionMapper,
+ IMaterializerRegistry materializerRegistry)
: IWorkflowDefinitionLinker
{
///
@@ -51,7 +53,7 @@ public class StaticWorkflowDefinitionLinker(
foreach (var item in list.Items)
{
- items.Add(new LinkedWorkflowDefinitionSummary
+ items.Add(new()
{
Links = GenerateLinksForSingleEntry(item.DefinitionId, item.IsReadonly),
Id = item.Id,
@@ -65,11 +67,12 @@ public class StaticWorkflowDefinitionLinker(
ProviderName = item.ProviderName,
MaterializerName = item.MaterializerName,
CreatedAt = item.CreatedAt,
- IsReadonly = item.IsReadonly
+ IsReadonly = item.IsReadonly,
+ IsMaterializerAvailable = materializerRegistry.IsMaterializerAvailable(item.MaterializerName)
});
}
- return new PagedListResponse
+ return new()
{
TotalCount = list.TotalCount,
Items = items,
@@ -126,13 +129,13 @@ public class StaticWorkflowDefinitionLinker(
if (!managementOptions.Value.IsReadOnlyMode)
{
- linksList.Add(new Link($"/bulk-actions/delete/workflow-definitions/by-definition-id", "bulk-delete-by-definition-id", "POST"));
- linksList.Add(new Link($"/bulk-actions/delete/workflow-definitions/by-id", "bulk-delete-by-id", "POST"));
- linksList.Add(new Link($"/bulk-actions/publish/workflow-definitions/by-definition-ids", "bulk-publish", "POST"));
- linksList.Add(new Link($"/bulk-actions/retract/workflow-definitions/by-definition-ids", "bulk-retract", "POST"));
- linksList.Add(new Link($"/workflow-definitions/import", "import", "POST"));
- linksList.Add(new Link($"/workflow-definitions/import-files", "import-files", "POST"));
- linksList.Add(new Link($"/workflow-definitions", "create", "POST"));
+ linksList.Add(new($"/bulk-actions/delete/workflow-definitions/by-definition-id", "bulk-delete-by-definition-id", "POST"));
+ linksList.Add(new($"/bulk-actions/delete/workflow-definitions/by-id", "bulk-delete-by-id", "POST"));
+ linksList.Add(new($"/bulk-actions/publish/workflow-definitions/by-definition-ids", "bulk-publish", "POST"));
+ linksList.Add(new($"/bulk-actions/retract/workflow-definitions/by-definition-ids", "bulk-retract", "POST"));
+ linksList.Add(new($"/workflow-definitions/import", "import", "POST"));
+ linksList.Add(new($"/workflow-definitions/import-files", "import-files", "POST"));
+ linksList.Add(new($"/workflow-definitions", "create", "POST"));
}
return linksList.ToArray();
@@ -153,11 +156,11 @@ public class StaticWorkflowDefinitionLinker(
if (!managementOptions.Value.IsReadOnlyMode && !definitionIsReadonly)
{
- links.Add(new Link($"/workflow-definitions/{definitionId}/publish", "publish", "POST"));
- links.Add(new Link($"/workflow-definitions/{definitionId}/retract", "retract", "POST"));
- links.Add(new Link($"/workflow-definitions/{definitionId}", "delete", "DELETE"));
- links.Add(new Link($"/workflow-definitions/{definitionId}/import", "import", "PUT"));
- links.Add(new Link($"/workflow-definitions/{definitionId}/update-references", "update-references", "POST"));
+ links.Add(new($"/workflow-definitions/{definitionId}/publish", "publish", "POST"));
+ links.Add(new($"/workflow-definitions/{definitionId}/retract", "retract", "POST"));
+ links.Add(new($"/workflow-definitions/{definitionId}", "delete", "DELETE"));
+ links.Add(new($"/workflow-definitions/{definitionId}/import", "import", "PUT"));
+ links.Add(new($"/workflow-definitions/{definitionId}/update-references", "update-references", "POST"));
}
return links.ToArray();
diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IMaterializerRegistry.cs b/src/modules/Elsa.Workflows.Management/Contracts/IMaterializerRegistry.cs
new file mode 100644
index 000000000..b314e7cfc
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Management/Contracts/IMaterializerRegistry.cs
@@ -0,0 +1,26 @@
+namespace Elsa.Workflows.Management;
+
+///
+/// Provides access to registered workflow materializers and their availability.
+///
+public interface IMaterializerRegistry
+{
+ ///
+ /// Gets all registered materializers.
+ ///
+ IEnumerable GetMaterializers();
+
+ ///
+ /// Gets a materializer by name.
+ ///
+ /// The name of the materializer.
+ /// The materializer, or null if not found.
+ IWorkflowMaterializer? GetMaterializer(string name);
+
+ ///
+ /// Checks if a materializer with the specified name is available.
+ ///
+ /// The name of the materializer.
+ /// True if the materializer is available, false otherwise.
+ bool IsMaterializerAvailable(string name);
+}
diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs
index 149591861..0ffa08dd3 100644
--- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs
+++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs
@@ -2,6 +2,7 @@ using Elsa.Common.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
+using Elsa.Workflows.Management.Models;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management;
@@ -60,4 +61,62 @@ public interface IWorkflowDefinitionService
/// Looks for all s that match the specified .
///
Task> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
+
+ ///
+ /// Attempts to retrieve a and its corresponding
+ /// based on the specified workflow definition ID and version options.
+ ///
+ /// The ID of the workflow definition to search for.
+ /// The versioning options that determine which version of the workflow definition to search for.
+ /// A token to monitor for cancellation requests.
+ ///
+ /// A containing the workflow graph and its associated definition,
+ /// or indicating that they do not exist.
+ ///
+ Task TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default);
+
+ ///
+ /// Attempts to locate a based on the specified workflow definition version ID.
+ ///
+ /// The unique identifier of the workflow definition version.
+ /// A token to monitor for cancellation requests.
+ ///
+ /// A containing the workflow graph and its associated definition,
+ /// or indicating that they do not exist.
+ ///
+ Task TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default);
+
+ ///
+ /// Attempts to retrieve the workflow graph corresponding to the specified .
+ ///
+ /// The handle identifying the workflow definition and its version.
+ /// A token to cancel the asynchronous operation.
+ ///
+ /// A containing the workflow graph and its associated definition,
+ /// or indicating that they do not exist.
+ ///
+ Task TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default);
+
+ ///
+ /// Attempts to retrieve a and its associated
+ /// for the specified workflow definition using the provided filter criteria.
+ ///
+ /// The criteria used to filter and locate the workflow definition and graph.
+ /// A token to monitor for cancellation requests.
+ ///
+ /// A containing the workflow graph and its associated definition,
+ /// or indicating that they do not exist.
+ ///
+ Task TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
+
+ ///
+ /// Attempts to find and retrieve a collection of objects based on the specified .
+ ///
+ /// The filter specifying criteria to identify the workflow definitions to search for.
+ /// A token to observe while waiting for the task to complete.
+ ///
+ /// A collection of instances, each containing a workflow graph and its associated definition
+ /// for every matching workflow definition. If no matching workflow definitions are found, the collection is empty.
+ ///
+ Task> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Management/Exceptions/WorkflowMaterializerNotFoundException.cs b/src/modules/Elsa.Workflows.Management/Exceptions/WorkflowMaterializerNotFoundException.cs
new file mode 100644
index 000000000..118ec869d
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Management/Exceptions/WorkflowMaterializerNotFoundException.cs
@@ -0,0 +1,6 @@
+namespace Elsa.Workflows.Management.Exceptions;
+
+public class WorkflowMaterializerNotFoundException(string materializerName, string message = "Materializer not found. The materializer may be disabled or not registered") : Exception(message)
+{
+ public string MaterializerName { get; } = materializerName;
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs
index 4740bf1e5..cd3da8fec 100644
--- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs
+++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs
@@ -255,6 +255,7 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module)
.AddScoped()
.AddScoped()
.AddScoped()
+ .AddScoped()
.AddScoped()
.AddScoped()
.AddScoped()
diff --git a/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs b/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs
index 1d02f2970..8dafcd254 100644
--- a/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs
+++ b/src/modules/Elsa.Workflows.Management/Models/TimestampFilter.cs
@@ -12,7 +12,7 @@ public class TimestampFilter
///
/// Gets or sets the column to filter by.
///
- public required string Column { get; set; } = default!;
+ public required string Column { get; set; } = null!;
///
/// Gets or sets the operator to use for filtering.
diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphFindResult.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphFindResult.cs
new file mode 100644
index 000000000..d3df0d968
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowGraphFindResult.cs
@@ -0,0 +1,10 @@
+using Elsa.Workflows.Management.Entities;
+using Elsa.Workflows.Models;
+
+namespace Elsa.Workflows.Management.Models;
+
+public record WorkflowGraphFindResult(WorkflowDefinition? WorkflowDefinition, WorkflowGraph? WorkflowGraph)
+{
+ public bool WorkflowDefinitionExists => WorkflowDefinition != null;
+ public bool WorkflowGraphExists => WorkflowGraph != null;
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs
index 059457ba8..578c1fe84 100644
--- a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs
+++ b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs
@@ -1,6 +1,7 @@
using Elsa.Common.Models;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
+using Elsa.Workflows.Management.Models;
using Elsa.Workflows.Models;
using JetBrains.Annotations;
using Microsoft.Extensions.Caching.Memory;
@@ -11,12 +12,16 @@ namespace Elsa.Workflows.Management.Services;
/// Decorates an with caching capabilities.
///
[UsedImplicitly]
-public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowDefinitionService
+public class CachingWorkflowDefinitionService(
+ IWorkflowDefinitionService decoratedService,
+ IWorkflowDefinitionCacheManager cacheManager,
+ IWorkflowDefinitionStore workflowDefinitionStore,
+ IMaterializerRegistry materializerRegistry) : IWorkflowDefinitionService
{
///
- public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
+ public Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
- return await decoratedService.MaterializeWorkflowAsync(definition, cancellationToken);
+ return decoratedService.MaterializeWorkflowAsync(definition, cancellationToken);
}
///
@@ -24,7 +29,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
{
var cacheKey = cacheManager.CreateWorkflowDefinitionVersionCacheKey(definitionId, versionOptions);
- return await GetFromCacheAsync(cacheKey,
+ return await FindFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken),
x => x.DefinitionId);
}
@@ -33,7 +38,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
public async Task FindWorkflowDefinitionAsync(string definitionVersionId, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowDefinitionVersionCacheKey(definitionVersionId);
- return await GetFromCacheAsync(cacheKey,
+ return await FindFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken),
x => x.DefinitionId);
}
@@ -41,8 +46,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
///
public Task FindWorkflowDefinitionAsync(WorkflowDefinitionHandle handle, CancellationToken cancellationToken = default)
{
- var filter = new WorkflowDefinitionFilter { DefinitionHandle = handle };
-
+ var filter = handle.ToFilter();
return FindWorkflowDefinitionAsync(filter, cancellationToken);
}
@@ -50,7 +54,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
public async Task FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowDefinitionFilterCacheKey(filter);
- return await GetFromCacheAsync(cacheKey,
+ return await FindFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowDefinitionAsync(filter, cancellationToken),
x => x.DefinitionId);
}
@@ -59,7 +63,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
public async Task FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionId, versionOptions);
- return await GetFromCacheAsync(cacheKey,
+ return await FindFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken),
x => x.Workflow.Identity.DefinitionId);
}
@@ -68,7 +72,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
public async Task FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionVersionId);
- return await GetFromCacheAsync(cacheKey,
+ return await FindFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowGraphAsync(definitionVersionId, cancellationToken),
x => x.Workflow.Identity.DefinitionId);
}
@@ -76,8 +80,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
///
public Task FindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default)
{
- var filter = new WorkflowDefinitionFilter { DefinitionHandle = definitionHandle };
-
+ var filter = definitionHandle.ToFilter();
return FindWorkflowGraphAsync(filter, cancellationToken);
}
@@ -85,11 +88,12 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
public async Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var cacheKey = cacheManager.CreateWorkflowFilterCacheKey(filter);
- return await GetFromCacheAsync(cacheKey,
+ return await FindFromCacheAsync(cacheKey,
() => decoratedService.FindWorkflowGraphAsync(filter, cancellationToken),
x => x.Workflow.Identity.DefinitionId);
}
+ ///
public async Task> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken);
@@ -101,26 +105,102 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
cacheKey,
async () => await MaterializeWorkflowAsync(workflowDefinition, cancellationToken),
wf => wf.Workflow.Identity.DefinitionId);
- workflowGraphs.Add(workflowGraph!);
+ workflowGraphs.Add(workflowGraph);
}
return workflowGraphs;
}
- private async Task GetFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc)
+ ///
+ public async Task TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
+ {
+ var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionId, versionOptions);
+ var result = await GetFromCacheAsync(cacheKey,
+ () => decoratedService.TryFindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken),
+ x => x.WorkflowDefinition?.DefinitionId);
+
+ return result;
+ }
+
+ ///
+ public async Task TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default)
+ {
+ var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionVersionId);
+ var result = await GetFromCacheAsync(cacheKey,
+ () => decoratedService.TryFindWorkflowGraphAsync(definitionVersionId, cancellationToken),
+ x => x.WorkflowDefinition?.DefinitionId);
+
+ return result;
+ }
+
+ ///
+ public Task TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default)
+ {
+ var filter = definitionHandle.ToFilter();
+ return TryFindWorkflowGraphAsync(filter, cancellationToken);
+ }
+
+ ///
+ public async Task TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
+ {
+ var cacheKey = cacheManager.CreateWorkflowFilterCacheKey(filter);
+ var result = await GetFromCacheAsync(cacheKey,
+ () => decoratedService.TryFindWorkflowGraphAsync(filter, cancellationToken),
+ x => x.WorkflowDefinition?.DefinitionId);
+
+ return result;
+ }
+
+ ///
+ public async Task> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
+ {
+ var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken);
+ var results = new List();
+ foreach (var workflowDefinition in workflowDefinitions)
+ {
+ if (!materializerRegistry.IsMaterializerAvailable(workflowDefinition.MaterializerName))
+ {
+ var unavailableResult = new WorkflowGraphFindResult(workflowDefinition, null);
+ results.Add(unavailableResult);
+ continue;
+ }
+
+ var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(workflowDefinition.Id);
+ var workflowGraph = await FindFromCacheAsync(
+ cacheKey,
+ async () => await MaterializeWorkflowAsync(workflowDefinition, cancellationToken),
+ wf => wf.Workflow.Identity.DefinitionId);
+
+ var result = new WorkflowGraphFindResult(workflowDefinition, workflowGraph);
+ results.Add(result);
+ }
+
+ return results;
+ }
+
+ private async Task FindFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc) where T : class
+ {
+ return await GetFromCacheAsync(
+ cacheKey,
+ getObjectFunc,
+ obj => obj != null ? getChangeTokenKeyFunc(obj) : null);
+ }
+
+ private async Task GetFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc)
{
var cache = cacheManager.Cache;
return await cache.GetOrCreateAsync(cacheKey, async entry =>
{
entry.SetAbsoluteExpiration(cache.CachingOptions.Value.CacheDuration);
var obj = await getObjectFunc();
-
- if (obj == null)
- return default;
-
var changeTokenKeyInput = getChangeTokenKeyFunc(obj);
- var changeTokenKey = cacheManager.CreateWorkflowDefinitionChangeTokenKey(changeTokenKeyInput);
- entry.AddExpirationToken(cache.GetToken(changeTokenKey));
+
+ if (changeTokenKeyInput != null)
+ {
+ var changeTokenKey = cacheManager.CreateWorkflowDefinitionChangeTokenKey(changeTokenKeyInput);
+ entry.AddExpirationToken(cache.GetToken(changeTokenKey));
+ }
+
return obj;
});
}
diff --git a/src/modules/Elsa.Workflows.Management/Services/MaterializerRegistry.cs b/src/modules/Elsa.Workflows.Management/Services/MaterializerRegistry.cs
new file mode 100644
index 000000000..f7239e881
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Management/Services/MaterializerRegistry.cs
@@ -0,0 +1,22 @@
+namespace Elsa.Workflows.Management.Services;
+
+///
+public class MaterializerRegistry(Func> materializers) : IMaterializerRegistry
+{
+ private readonly Lazy> _materializers = new(() => materializers().ToArray());
+
+ ///
+ public IEnumerable GetMaterializers() => _materializers.Value;
+
+ ///
+ public IWorkflowMaterializer? GetMaterializer(string name)
+ {
+ return _materializers.Value.FirstOrDefault(x => x.Name == name);
+ }
+
+ ///
+ public bool IsMaterializerAvailable(string name)
+ {
+ return _materializers.Value.Any(x => x.Name == name);
+ }
+}
diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs
index 104ea4e8c..a91107ee7 100644
--- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs
+++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs
@@ -1,7 +1,10 @@
using Elsa.Common.Models;
using Elsa.Workflows.Management.Entities;
+using Elsa.Workflows.Management.Exceptions;
using Elsa.Workflows.Management.Filters;
+using Elsa.Workflows.Management.Models;
using Elsa.Workflows.Models;
+using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Management.Services;
@@ -9,17 +12,17 @@ namespace Elsa.Workflows.Management.Services;
public class WorkflowDefinitionService(
IWorkflowDefinitionStore workflowDefinitionStore,
IWorkflowGraphBuilder workflowGraphBuilder,
- Func> materializers)
+ IMaterializerRegistry materializerRegistry,
+ ILogger logger)
: IWorkflowDefinitionService
{
///
public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
- var workflowMaterializers = materializers();
- var materializer = workflowMaterializers.FirstOrDefault(x => x.Name == definition.MaterializerName);
+ var materializer = materializerRegistry.GetMaterializer(definition.MaterializerName);
if (materializer == null)
- throw new("Provider not found");
+ throw new WorkflowMaterializerNotFoundException(definition.MaterializerName);
var workflow = await materializer.MaterializeAsync(definition, cancellationToken);
return await workflowGraphBuilder.BuildAsync(workflow, cancellationToken);
@@ -56,44 +59,28 @@ public class WorkflowDefinitionService(
public async Task FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken);
-
- if (definition == null)
- return null;
-
- return await MaterializeWorkflowAsync(definition, cancellationToken);
+ return await TryMaterializeWorkflowAsync(definition, cancellationToken);
}
///
public async Task FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken);
-
- if (definition == null)
- return null;
-
- return await MaterializeWorkflowAsync(definition, cancellationToken);
+ return await TryMaterializeWorkflowAsync(definition, cancellationToken);
}
///
public async Task FindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(definitionHandle, cancellationToken);
-
- if (definition == null)
- return null;
-
- return await MaterializeWorkflowAsync(definition, cancellationToken);
+ return await TryMaterializeWorkflowAsync(definition, cancellationToken);
}
///
public async Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken);
-
- if (definition == null)
- return null;
-
- return await MaterializeWorkflowAsync(definition, cancellationToken);
+ return await TryMaterializeWorkflowAsync(definition, cancellationToken);
}
///
@@ -109,4 +96,81 @@ public class WorkflowDefinitionService(
return workflowGraphs;
}
+
+ public async Task TryFindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
+ {
+ var definition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken);
+ return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
+ }
+
+ public async Task TryFindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default)
+ {
+ var definition = await FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken);
+ return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
+ }
+
+ public async Task TryFindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default)
+ {
+ var definition = await FindWorkflowDefinitionAsync(definitionHandle, cancellationToken);
+ return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
+ }
+
+ public async Task TryFindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
+ {
+ var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken);
+ return await MaterializeWorkflowGraphFindResultAsync(definition, cancellationToken);
+ }
+
+ public async Task> TryFindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
+ {
+ var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken);
+ var results = new List();
+ foreach (var workflowDefinition in workflowDefinitions)
+ {
+ var result = await MaterializeWorkflowGraphFindResultAsync(workflowDefinition, cancellationToken);
+ results.Add(result);
+ }
+
+ return results;
+ }
+
+ ///
+ /// Attempts to materialize a workflow graph from the given workflow definition if a suitable materializer is available.
+ ///
+ /// The workflow definition to materialize. Can be null.
+ /// A token to monitor for cancellation requests.
+ ///
+ /// A if materialization is successful; otherwise, null.
+ ///
+ private async Task TryMaterializeWorkflowAsync(WorkflowDefinition? definition, CancellationToken cancellationToken)
+ {
+ if (definition == null)
+ return null;
+
+ if (materializerRegistry.IsMaterializerAvailable(definition.MaterializerName))
+ return await MaterializeWorkflowAsync(definition, cancellationToken);
+
+ logger.LogWarning("Materializer '{MaterializerName}' not found. The workflow definition will not be materialized.", definition.MaterializerName);
+ return null;
+ }
+
+ ///
+ /// Attempts to materialize a workflow graph from the given workflow definition.
+ ///
+ /// The workflow definition to materialize the graph for. May be null.
+ /// A token to observe while waiting for the task to complete.
+ /// A result containing the materialized workflow graph and its corresponding workflow definition, if successfully materialized; otherwise, returns the definition with a null graph.
+ private async Task MaterializeWorkflowGraphFindResultAsync(WorkflowDefinition? definition, CancellationToken cancellationToken)
+ {
+ if (definition == null)
+ return new(null, null);
+
+ if (materializerRegistry.IsMaterializerAvailable(definition.MaterializerName))
+ {
+ var graph = await MaterializeWorkflowAsync(definition, cancellationToken);
+ return new(definition, graph);
+ }
+
+ return new(definition, null);
+ }
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs b/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs
index 796a55755..22132cd2e 100644
--- a/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs
+++ b/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs
@@ -140,7 +140,7 @@ public class CachingWorkflowDefinitionStore(IWorkflowDefinitionStore decoratedSt
var tenantId = tenantAccessor.Tenant?.Id;
var tenantIdPrefix = !string.IsNullOrEmpty(tenantId) ? $"{tenantId}:" : string.Empty;
var internalKey = $"{tenantIdPrefix}{typeof(T).Name}:{key}";
- return await cacheManager.GetOrCreateAsync(internalKey, async entry =>
+ return await cacheManager.FindOrCreateAsync(internalKey, async entry =>
{
var invalidationRequestToken = cacheManager.GetToken(CacheInvalidationTokenKey);
entry.AddExpirationToken(invalidationRequestToken);
diff --git a/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowDefinitionNotFoundException.cs b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowDefinitionNotFoundException.cs
new file mode 100644
index 000000000..de7ac14d2
--- /dev/null
+++ b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowDefinitionNotFoundException.cs
@@ -0,0 +1,8 @@
+using Elsa.Workflows.Models;
+
+namespace Elsa.Workflows.Runtime.Exceptions;
+
+public class WorkflowDefinitionNotFoundException(string message, WorkflowDefinitionHandle workflowDefinitionHandle) : Exception(message)
+{
+ public WorkflowDefinitionHandle WorkflowDefinitionHandle { get; } = workflowDefinitionHandle;
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs
index 1d0f302fe..bee3b00a0 100644
--- a/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Exceptions/WorkflowGraphNotFoundException.cs
@@ -2,6 +2,7 @@ using Elsa.Workflows.Models;
namespace Elsa.Workflows.Runtime.Exceptions;
+[Obsolete("Use WorkflowDefinitionNotFoundException instead.")]
public class WorkflowGraphNotFoundException(string message, WorkflowDefinitionHandle workflowDefinitionHandle) : Exception(message)
{
public WorkflowDefinitionHandle WorkflowDefinitionHandle { get; } = workflowDefinitionHandle;
diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs
index a07d3823f..761f8396f 100644
--- a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowCancellationDispatcher.cs
@@ -1,21 +1,28 @@
-using Elsa.Mediator;
-using Elsa.Mediator.Contracts;
-using Elsa.Workflows.Runtime.Commands;
-using Elsa.Workflows.Runtime.Requests;
-using Elsa.Workflows.Runtime.Responses;
-
-namespace Elsa.Workflows.Runtime;
-
-///
-/// Dispatches workflow cancellation requests to a local background worker.
-///
-public class BackgroundWorkflowCancellationDispatcher(ICommandSender commandSender) : IWorkflowCancellationDispatcher
-{
- ///
- public async Task DispatchAsync(DispatchCancelWorkflowRequest request, CancellationToken cancellationToken = default)
- {
- var command = new CancelWorkflowsCommand(request);
- await commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken);
- return new DispatchCancelWorkflowsResponse();
- }
+using Elsa.Common.Multitenancy;
+using Elsa.Mediator;
+using Elsa.Mediator.Contracts;
+using Elsa.Tenants.Mediator;
+using Elsa.Workflows.Runtime.Commands;
+using Elsa.Workflows.Runtime.Requests;
+using Elsa.Workflows.Runtime.Responses;
+
+namespace Elsa.Workflows.Runtime;
+
+///
+/// Dispatches workflow cancellation requests to a local background worker.
+///
+public class BackgroundWorkflowCancellationDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IWorkflowCancellationDispatcher
+{
+ ///
+ public async Task DispatchAsync(DispatchCancelWorkflowRequest request, CancellationToken cancellationToken = default)
+ {
+ var command = new CancelWorkflowsCommand(request);
+ await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken);
+ return new();
+ }
+
+ private IDictionary