From 21c54f583eeea2b83aef536421bd8ffb01243e60 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 15 Apr 2024 10:06:27 +0200 Subject: [PATCH] Fix WorkflowActivity to Use Cached Workflow Definitions for Consistent Behavior (#5223) * Replace IServiceScopeFactory with IServiceProvider in WorkflowRunner Unused dependencies were removed from the workflow runner service. The IServiceScopeFactory was replaced with IServiceProvider to better handle the creation and deletion of service scopes, resulting in cleaner code with less manual scope management. Microsoft.Extensions.DependencyInjection and System.Diagnostics.CodeAnalysis were removed as they were no longer necessary. * Refactor WorkflowDefinitionActivity to use WorkflowDefinitionService The WorkflowDefinitionActivity class has been refactored to make use of the WorkflowDefinitionService instead of the WorkflowDefinitionStore. This fixes #5222 by ensuring the same activity instances are used in the graph model of the workflow execution context. * Add workflow filtering and caching functionality Added methods to `WorkflowDefinitionService` to find workflow definitions and workflows using filter criteria. A key generation method for caching filtered workflows was also added to `WorkflowDefinitionCacheManager`. The implementation includes generating a hash of the filter parameters and using this hash as a cache key, providing efficient caching functionality for filtered searches. * Refactor TriggerIndexer to handle only ITrigger activities The code in TriggerIndexer has been refactored to deal specifically with ITrigger activities, streamlining its behavior. Removed code related to handling non-ITrigger activities and simplified the workflow creation process. The extraction of "startable" nodes now directly filters and casts to ITrigger, reducing complexity and increasing readability. * Update caching service to support filter-based search The CachingWorkflowDefinitionService has been updated to support workflow definition and workflow search based on filter criteria. The update also includes change of class scope from public to internal. Further, it resolves the missing reference by switching from Elsa.Caching.Contracts to Elsa.Caching. * Optimize Elsa project imports and use explicit cache variable names This commit removes superfluous import references, relocates the 'IChangeTokenSignaler' contract into the 'Elsa.Caching' namespace, and replaces ambiguous 'cache' variable names with more explicit 'memoryCache' across several files. Additionally, new package references have been added and access modifiers have been changed to improve encapsulation. Cleanup enhances readability and maintainability of the codebase. * Add .DotSettings file to Elsa.Caching module A new .DotSettings file has been added to the Elsa.Caching module. This file is used for namespace configuration, specifically to skip the "contracts" folder in code inspections. * Update workflow interfaces to support filter queries The update extends `IWorkflowDefinitionCacheManager` and `IWorkflowDefinitionService` interfaces. Functions are added to allow creating filter cache keys and finding workflow definitions and workflows using a new `WorkflowDefinitionFilter`. This enhances querying flexibility by enabling filtered searches. * Add WorkflowDefinitionVersionId to WorkflowTriggerEqualityComparer A new property, WorkflowDefinitionVersionId, has been added to the object being serialized in WorkflowTriggerEqualityComparer. This change allows for a more accurate comparison between workflow triggers, considering not just the workflow definition ID but also its version. * Update service registration types in WorkflowsFeature Changed the registration type for both IHasher and IBookmarkHasher services from Scoped to Singleton in the workflows feature configuration. This alteration aims to improve application performance and manage service lifetimes more efficiently. * Remove Open.Linq.AsyncExtensions dependency The Open.Linq.AsyncExtensions package reference was removed across the project. The usage within the CachingWorkflowDefinitionStore was updated accordingly to maintain functionality. * Move System.Linq.Dynamic.Core package reference The System.Linq.Dynamic.Core package reference was moved from the Directory.Build.props file to the Elsa.Workflows.Management.csproj file. This change reflects the specific dependency of the Elsa.Workflows.Management module on System.Linq.Dynamic.Core, without impacting other modules. * Implement caching for HTTP workflows This update introduces caching mechanisms for HTTP workflows, which significantly improves their performance. The changes involve creating a `CacheManager` and `CachingHttpWorkflowLookupService`, and modifying some existing components to use the new caching mechanism. Additionally, the `HttpWorkflowsCacheManager` was renamed to `HttpWorkflowsCacheInvalidationManager` to better reflect its role. * Refactor cache management across modules This commit refactor the cache management across various modules. The 'ICacheManager' interface now includes methods for triggering and getting change tokens, and the 'HttpWorkflowsCacheInvalidationManager' has been renamed to 'HttpWorkflowsCacheManager'. The caching functionality in 'WorkflowDefinitionService' and other similar services have been updated to use these new methods, improving consistency and maintainability. * Enable caching in Elsa.Server.Web The "useCaching" variable has been set to true to enable caching. Simultaneously, the method name "UseCachingStores" has been refactored to "UseCache". Conditional statements have been added to check the "useCaching" variable before invoking caching. * Rename method UseCaching to UseCache In the Elsa.Server.Web and Elsa.Http project files, the method UseCaching has been renamed to UseCache. This modification is aimed at bridging naming inconsistencies and maintaining naming standards across the application. * Update HttpCacheFeature class description The class summary for HttpCacheFeature has been revised. Originally, it stated that the class was used for installing services related to HTTP services and activities, but it actually focuses more on HTTP caching. * Remove unused Configure method from HttpCacheFeature The Configure method in HttpCacheFeature was found to be redundant as it wasn't doing any significant work or contributing to any functionality. It has therefore been removed to clean up the code and avoid confusion. * Add 'bug/*' to workflow triggers This commit includes 'bug/*' to the list of triggers in our GitHub Actions workflow. Now, any push or pull request under a 'bug/*' branch will trigger the workflow. --- .github/workflows/packages.yml | 1 + src/bundles/Elsa.Server.Web/Program.cs | 15 ++-- .../MassTransitChangeTokenSignalPublisher.cs | 3 +- .../Features/DistributedCacheFeature.cs | 3 +- .../DistributedChangeTokenSignaler.cs | 1 - .../NoopChangeTokenSignalPublisher.cs | 1 - .../Elsa.Caching/Contracts/ICacheManager.cs | 32 ++++++++ .../Contracts/IChangeTokenSignaler.cs | 2 +- .../Elsa.Caching.csproj.DotSettings | 2 + .../Features/MemoryCacheFeature.cs | 4 +- .../Elsa.Caching/Services/CacheManager.cs | 31 ++++++++ .../Services/ChangeTokenSignaler.cs | 1 - .../Contracts/IHttpWorkflowLookupService.cs | 14 ++++ .../Contracts/IHttpWorkflowsCacheManager.cs | 19 +++-- .../Elsa.Http/Extensions/ModuleExtensions.cs | 10 ++- .../Elsa.Http/Features/HttpCacheFeature.cs | 32 ++++++++ src/modules/Elsa.Http/Features/HttpFeature.cs | 6 +- .../Middleware/HttpWorkflowsMiddleware.cs | 12 +-- .../Models/HttpWorkflowLookupResult.cs | 6 ++ .../CachingHttpWorkflowLookupService.cs | 40 ++++++++++ .../Services/HttpWorkflowLookupService.cs | 50 ++++++++++++ .../Services/HttpWorkflowsCacheManager.cs | 76 +++---------------- .../Features/WorkflowsFeature.cs | 4 +- .../Services/WorkflowRunner.cs | 22 ++---- .../WorkflowDefinitionActivity.cs | 28 +++---- .../IWorkflowDefinitionCacheManager.cs | 11 +++ .../Contracts/IWorkflowDefinitionService.cs | 11 +++ .../Elsa.Workflows.Management.csproj | 3 +- .../Extensions/ModuleExtensions.cs | 2 +- .../EvictWorkflowDefinitionServiceCache.cs | 2 +- .../CachingWorkflowDefinitionService.cs | 36 ++++++--- .../WorkflowDefinitionCacheManager.cs | 25 +++++- .../Services/WorkflowDefinitionService.cs | 17 +++++ .../Stores/CachingWorkflowDefinitionStore.cs | 27 +++---- .../WorkflowTriggerEqualityComparer.cs | 1 + .../Extensions/ModuleExtensions.cs | 2 +- .../Services/TriggerIndexer.cs | 42 ++-------- .../Stores/CachingBookmarkStore.cs | 23 ++---- .../Stores/CachingTriggerStore.cs | 25 +++--- 39 files changed, 407 insertions(+), 235 deletions(-) create mode 100644 src/modules/Elsa.Caching/Contracts/ICacheManager.cs create mode 100644 src/modules/Elsa.Caching/Elsa.Caching.csproj.DotSettings create mode 100644 src/modules/Elsa.Caching/Services/CacheManager.cs create mode 100644 src/modules/Elsa.Http/Contracts/IHttpWorkflowLookupService.cs create mode 100644 src/modules/Elsa.Http/Features/HttpCacheFeature.cs create mode 100644 src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs create mode 100644 src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs create mode 100644 src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index 2bc10c7dd..05a7ed157 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -6,6 +6,7 @@ on: - 'main' - 'feature/*' - 'issue/*' + - 'bug/*' - 'enhancement/*' - 'patch/*' - 'fix/*' diff --git a/src/bundles/Elsa.Server.Web/Program.cs b/src/bundles/Elsa.Server.Web/Program.cs index ef688eeb9..52bac227a 100644 --- a/src/bundles/Elsa.Server.Web/Program.cs +++ b/src/bundles/Elsa.Server.Web/Program.cs @@ -42,7 +42,7 @@ const bool useMassTransit = true; const bool useZipCompression = true; const bool runEFCoreMigrations = true; const bool useMemoryStores = false; -const bool useCachingStores = true; +const bool useCaching = true; const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit; const MassTransitBroker useMassTransitBroker = MassTransitBroker.Memory; @@ -153,9 +153,9 @@ services if (useMassTransit) management.UseMassTransitDispatcher(); - if (useCachingStores) - management.UseCachingStores(); - + if (useCaching) + management.UseCache(); + management.SetDefaultLogPersistenceMode(LogPersistenceMode.Default); }) .UseWorkflowRuntime(runtime => @@ -203,8 +203,8 @@ services runtime.WorkflowInboxStore = sp => sp.GetRequiredService(); } - if (useCachingStores) - runtime.UseCachingStores(); + if (useCaching) + runtime.UseCache(); runtime.DistributedLockProvider = _ => { @@ -268,6 +268,9 @@ services .UseHttp(http => { http.ConfigureHttpOptions = options => configuration.GetSection("Http").Bind(options); + + if (useCaching) + http.UseCache(); }) .UseEmail(email => email.ConfigureOptions = options => configuration.GetSection("Smtp").Bind(options)) .UseAlterations(alterations => diff --git a/src/modules/Elsa.Caching.Distributed.MassTransit/Services/MassTransitChangeTokenSignalPublisher.cs b/src/modules/Elsa.Caching.Distributed.MassTransit/Services/MassTransitChangeTokenSignalPublisher.cs index d159313bd..cf60fb2ea 100644 --- a/src/modules/Elsa.Caching.Distributed.MassTransit/Services/MassTransitChangeTokenSignalPublisher.cs +++ b/src/modules/Elsa.Caching.Distributed.MassTransit/Services/MassTransitChangeTokenSignalPublisher.cs @@ -1,5 +1,4 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Distributed.Contracts; +using Elsa.Caching.Distributed.Contracts; using Elsa.Caching.Distributed.MassTransit.Messages; using MassTransit; diff --git a/src/modules/Elsa.Caching.Distributed/Features/DistributedCacheFeature.cs b/src/modules/Elsa.Caching.Distributed/Features/DistributedCacheFeature.cs index b63a62d0c..d519a728a 100644 --- a/src/modules/Elsa.Caching.Distributed/Features/DistributedCacheFeature.cs +++ b/src/modules/Elsa.Caching.Distributed/Features/DistributedCacheFeature.cs @@ -1,5 +1,4 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Distributed.Contracts; +using Elsa.Caching.Distributed.Contracts; using Elsa.Caching.Distributed.Services; using Elsa.Caching.Features; using Elsa.Caching.Options; diff --git a/src/modules/Elsa.Caching.Distributed/Services/DistributedChangeTokenSignaler.cs b/src/modules/Elsa.Caching.Distributed/Services/DistributedChangeTokenSignaler.cs index b9a1dd5f0..98de45eeb 100644 --- a/src/modules/Elsa.Caching.Distributed/Services/DistributedChangeTokenSignaler.cs +++ b/src/modules/Elsa.Caching.Distributed/Services/DistributedChangeTokenSignaler.cs @@ -1,4 +1,3 @@ -using Elsa.Caching.Contracts; using Elsa.Caching.Distributed.Contracts; using JetBrains.Annotations; using Microsoft.Extensions.Primitives; diff --git a/src/modules/Elsa.Caching.Distributed/Services/NoopChangeTokenSignalPublisher.cs b/src/modules/Elsa.Caching.Distributed/Services/NoopChangeTokenSignalPublisher.cs index ae5ada6f5..1eff9f288 100644 --- a/src/modules/Elsa.Caching.Distributed/Services/NoopChangeTokenSignalPublisher.cs +++ b/src/modules/Elsa.Caching.Distributed/Services/NoopChangeTokenSignalPublisher.cs @@ -1,4 +1,3 @@ -using Elsa.Caching.Contracts; using Elsa.Caching.Distributed.Contracts; namespace Elsa.Caching.Distributed.Services; diff --git a/src/modules/Elsa.Caching/Contracts/ICacheManager.cs b/src/modules/Elsa.Caching/Contracts/ICacheManager.cs new file mode 100644 index 000000000..6e93fe8c8 --- /dev/null +++ b/src/modules/Elsa.Caching/Contracts/ICacheManager.cs @@ -0,0 +1,32 @@ +using Elsa.Caching.Options; +using Microsoft.Extensions.Caching.Memory; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Primitives; + +namespace Elsa.Caching; + +/// +/// A thin wrapper around , allowing for centralized handling of cache entries. +/// +public interface ICacheManager +{ + /// + /// Provides options for configuring caching. + /// + IOptions CachingOptions { get; } + + /// + /// Gets a change token for the specified key. + /// + IChangeToken GetToken(string key); + + /// + /// Triggers the change token for the specified key. + /// + ValueTask TriggerTokenAsync(string key, CancellationToken cancellationToken = default); + + /// + /// Gets an item from the cache, or creates it if it doesn't exist. + /// + Task GetOrCreateAsync(object key, Func> factory); +} \ No newline at end of file diff --git a/src/modules/Elsa.Caching/Contracts/IChangeTokenSignaler.cs b/src/modules/Elsa.Caching/Contracts/IChangeTokenSignaler.cs index a6f27ea0c..d8ffe31b4 100644 --- a/src/modules/Elsa.Caching/Contracts/IChangeTokenSignaler.cs +++ b/src/modules/Elsa.Caching/Contracts/IChangeTokenSignaler.cs @@ -1,6 +1,6 @@ using Microsoft.Extensions.Primitives; -namespace Elsa.Caching.Contracts; +namespace Elsa.Caching; /// /// Provides change tokens for memory caches, allowing code to evict cache entries by triggering a signal. diff --git a/src/modules/Elsa.Caching/Elsa.Caching.csproj.DotSettings b/src/modules/Elsa.Caching/Elsa.Caching.csproj.DotSettings new file mode 100644 index 000000000..127787e89 --- /dev/null +++ b/src/modules/Elsa.Caching/Elsa.Caching.csproj.DotSettings @@ -0,0 +1,2 @@ + + True \ No newline at end of file diff --git a/src/modules/Elsa.Caching/Features/MemoryCacheFeature.cs b/src/modules/Elsa.Caching/Features/MemoryCacheFeature.cs index a96c98295..cc0eba6d5 100644 --- a/src/modules/Elsa.Caching/Features/MemoryCacheFeature.cs +++ b/src/modules/Elsa.Caching/Features/MemoryCacheFeature.cs @@ -1,5 +1,4 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Options; +using Elsa.Caching.Options; using Elsa.Caching.Services; using Elsa.Features.Abstractions; using Elsa.Features.Services; @@ -22,6 +21,7 @@ public class MemoryCacheFeature(IModule module) : FeatureBase(module) { Services.Configure(CachingOptions); Services.AddMemoryCache(); + Services.AddSingleton(); Services.AddSingleton(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Caching/Services/CacheManager.cs b/src/modules/Elsa.Caching/Services/CacheManager.cs new file mode 100644 index 000000000..6db501327 --- /dev/null +++ b/src/modules/Elsa.Caching/Services/CacheManager.cs @@ -0,0 +1,31 @@ +using Elsa.Caching.Options; +using Microsoft.Extensions.Caching.Memory; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Primitives; + +namespace Elsa.Caching.Services; + +/// +public class CacheManager(IMemoryCache memoryCache, IChangeTokenSignaler changeTokenSignaler, IOptions options) : ICacheManager +{ + /// + public IOptions CachingOptions => options; + + /// + public IChangeToken GetToken(string key) + { + return changeTokenSignaler.GetToken(key); + } + + /// + public ValueTask TriggerTokenAsync(string key, CancellationToken cancellationToken = default) + { + return changeTokenSignaler.TriggerTokenAsync(key, cancellationToken); + } + + /// + public async Task GetOrCreateAsync(object key, Func> factory) + { + return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry)); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Caching/Services/ChangeTokenSignaler.cs b/src/modules/Elsa.Caching/Services/ChangeTokenSignaler.cs index ef34c311f..0471da10c 100644 --- a/src/modules/Elsa.Caching/Services/ChangeTokenSignaler.cs +++ b/src/modules/Elsa.Caching/Services/ChangeTokenSignaler.cs @@ -1,5 +1,4 @@ using System.Collections.Concurrent; -using Elsa.Caching.Contracts; using Microsoft.Extensions.Primitives; namespace Elsa.Caching.Services; diff --git a/src/modules/Elsa.Http/Contracts/IHttpWorkflowLookupService.cs b/src/modules/Elsa.Http/Contracts/IHttpWorkflowLookupService.cs new file mode 100644 index 000000000..e6849efd1 --- /dev/null +++ b/src/modules/Elsa.Http/Contracts/IHttpWorkflowLookupService.cs @@ -0,0 +1,14 @@ +using Elsa.Http.Models; + +namespace Elsa.Http.Contracts; + +/// +/// Represents a service that looks up HTTP workflows and triggers. +/// +public interface IHttpWorkflowLookupService +{ + /// + /// Finds a workflow and trigger by bookmark hash. + /// + Task FindWorkflowAsync(string bookmarkHash, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Http/Contracts/IHttpWorkflowsCacheManager.cs b/src/modules/Elsa.Http/Contracts/IHttpWorkflowsCacheManager.cs index c6ed04e0b..25f0e6fa8 100644 --- a/src/modules/Elsa.Http/Contracts/IHttpWorkflowsCacheManager.cs +++ b/src/modules/Elsa.Http/Contracts/IHttpWorkflowsCacheManager.cs @@ -1,5 +1,4 @@ -using Elsa.Workflows.Activities; -using Elsa.Workflows.Runtime.Entities; +using Elsa.Caching; namespace Elsa.Http.Contracts; @@ -9,17 +8,27 @@ namespace Elsa.Http.Contracts; public interface IHttpWorkflowsCacheManager { /// - /// Finds a cached entry by bookmark hash. + /// Gets the cache manager. /// - Task<(Workflow? Workflow, ICollection Triggers)?> FindWorkflowAsync(string bookmarkHash, CancellationToken cancellationToken = default); + public ICacheManager Cache { get; } /// /// Evicts a cached entry by its definition ID. /// Task EvictWorkflowAsync(string workflowDefinitionId, CancellationToken cancellationToken = default); - + /// /// Evicts a cached entry by its bookmark hash. /// Task EvictTriggerAsync(string bookmarkHash, CancellationToken cancellationToken = default); + + /// + /// Gets the key for a workflow change token. + /// + string GetWorkflowChangeTokenKey(string workflowDefinitionId); + + /// + /// Gets the key for a trigger change token. + /// + string GetTriggerChangeTokenKey(string bookmarkHash); } \ No newline at end of file diff --git a/src/modules/Elsa.Http/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Http/Extensions/ModuleExtensions.cs index 20f93626c..7898a489c 100644 --- a/src/modules/Elsa.Http/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.Http/Extensions/ModuleExtensions.cs @@ -15,7 +15,15 @@ public static class ModuleExtensions public static IModule UseHttp(this IModule module, Action? configure = default) { module.Configure(configure); - module.Use(); return module; } + + /// + /// Install the feature to speed up HTTP workflows. Like, a lot. + /// + public static HttpFeature UseCache(this HttpFeature feature, Action? configure = default) + { + feature.Module.Configure(configure); + return feature; + } } \ No newline at end of file diff --git a/src/modules/Elsa.Http/Features/HttpCacheFeature.cs b/src/modules/Elsa.Http/Features/HttpCacheFeature.cs new file mode 100644 index 000000000..1f54df1da --- /dev/null +++ b/src/modules/Elsa.Http/Features/HttpCacheFeature.cs @@ -0,0 +1,32 @@ +using Elsa.Features.Abstractions; +using Elsa.Features.Attributes; +using Elsa.Features.Services; +using Elsa.Http.Contracts; +using Elsa.Http.Handlers; +using Elsa.Http.Services; +using Elsa.Workflows.Management.Features; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Http.Features; + +/// +/// Installs services related to HTTP caching. +/// +[DependsOn(typeof(HttpFeature))] +[DependsOn(typeof(CachingWorkflowDefinitionsFeature))] +public class HttpCacheFeature : FeatureBase +{ + /// + public HttpCacheFeature(IModule module) : base(module) + { + } + + /// + public override void Apply() + { + Services + .AddSingleton() + .Decorate() + .AddNotificationHandler(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Http/Features/HttpFeature.cs b/src/modules/Elsa.Http/Features/HttpFeature.cs index 6483e2cc6..4f326cc32 100644 --- a/src/modules/Elsa.Http/Features/HttpFeature.cs +++ b/src/modules/Elsa.Http/Features/HttpFeature.cs @@ -1,4 +1,3 @@ -using Elsa.Caching.Features; using Elsa.Expressions.Options; using Elsa.Extensions; using Elsa.Features.Abstractions; @@ -32,7 +31,7 @@ namespace Elsa.Http.Features; /// /// Installs services related to HTTP services and activities. /// -[DependsOn(typeof(MemoryCacheFeature))] +[DependsOn(typeof(HttpJavaScriptFeature))] public class HttpFeature : FeatureBase { /// @@ -151,14 +150,13 @@ public class HttpFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() - .AddScoped() + .AddScoped() .AddScoped(ContentTypeProvider) .AddHttpContextAccessor() // Handlers. .AddRequestHandler() .AddNotificationHandler() - .AddNotificationHandler() // Content parsers. .AddSingleton() diff --git a/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs b/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs index 4bf0e9812..a6db8b954 100644 --- a/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs +++ b/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs @@ -67,13 +67,13 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions(); + var httpWorkflowLookupService = serviceProvider.GetRequiredService(); var bookmarkHash = ComputeBookmarkHash(serviceProvider, matchingPath, method); - var cachedWorkflowAndTriggers = await httpEndpointCacheManager.FindWorkflowAsync(bookmarkHash, cancellationToken); + var lookupResult = await httpWorkflowLookupService.FindWorkflowAsync(bookmarkHash, cancellationToken); - if (cachedWorkflowAndTriggers != null) + if (lookupResult != null) { - var triggers = cachedWorkflowAndTriggers.Value.Triggers; + var triggers = lookupResult.Triggers; if (triggers?.Count > 1) { @@ -87,7 +87,7 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions { var correlationId = await GetCorrelationIdAsync(serviceProvider, httpContext, ct); diff --git a/src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs b/src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs new file mode 100644 index 000000000..2df870fa0 --- /dev/null +++ b/src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs @@ -0,0 +1,6 @@ +using Elsa.Workflows.Activities; +using Elsa.Workflows.Runtime.Entities; + +namespace Elsa.Http.Models; + +public record HttpWorkflowLookupResult(Workflow? Workflow, ICollection Triggers); \ No newline at end of file diff --git a/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs new file mode 100644 index 000000000..1f0061724 --- /dev/null +++ b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs @@ -0,0 +1,40 @@ +using Elsa.Http.Contracts; +using Elsa.Http.Models; +using JetBrains.Annotations; +using Microsoft.Extensions.Caching.Memory; + +namespace Elsa.Http.Services; + +/// +/// Represents a caching implementation of the IHttpWorkflowLookupService that retrieves workflows using HTTP workflow lookup service and caches the results. +/// +[UsedImplicitly] +public class CachingHttpWorkflowLookupService( + IHttpWorkflowLookupService decoratedService, + IHttpWorkflowsCacheManager cacheManager) : IHttpWorkflowLookupService +{ + /// + public async Task FindWorkflowAsync(string bookmarkHash, CancellationToken cancellationToken = default) + { + var key = $"http-workflow:{bookmarkHash}"; + var cache = cacheManager.Cache; + return await cache.GetOrCreateAsync(key, async entry => + { + var cachingOptions = cache.CachingOptions.Value; + entry.SetSlidingExpiration(cachingOptions.CacheDuration); + entry.AddExpirationToken(cache.GetToken(cacheManager.GetTriggerChangeTokenKey(bookmarkHash))); + + var result = await decoratedService.FindWorkflowAsync(bookmarkHash, cancellationToken); + + if (result == null) + return null; + + var workflow = result.Workflow!; + var changeTokenKey = cacheManager.GetWorkflowChangeTokenKey(workflow.Identity.DefinitionId); + var changeToken = cache.GetToken(changeTokenKey); + entry.AddExpirationToken(changeToken); + + return result; + }); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs b/src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs new file mode 100644 index 000000000..2eb860298 --- /dev/null +++ b/src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs @@ -0,0 +1,50 @@ +using Elsa.Http.Contracts; +using Elsa.Http.Models; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Entities; +using Elsa.Workflows.Runtime.Filters; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Http.Services; + +/// +public class HttpWorkflowLookupService(ITriggerStore triggerStore, IWorkflowDefinitionService workflowDefinitionService) : IHttpWorkflowLookupService +{ + /// + public async Task FindWorkflowAsync(string bookmarkHash, CancellationToken cancellationToken = default) + { + var triggers = await FindTriggersAsync(bookmarkHash, cancellationToken).ToList(); + + if (triggers.Count > 1) + return new(null, triggers); + + var trigger = triggers.SingleOrDefault(); + + if (trigger == null) + return default; + + var workflow = await FindWorkflowAsync(trigger, cancellationToken); + + if (workflow == null) + return default; + + return new(workflow, triggers); + } + + private async Task> FindTriggersAsync(string bookmarkHash, CancellationToken cancellationToken) + { + var triggerFilter = new TriggerFilter + { + Hash = bookmarkHash + }; + return await triggerStore.FindManyAsync(triggerFilter, cancellationToken); + } + + private async Task FindWorkflowAsync(StoredTrigger trigger, CancellationToken cancellationToken) + { + var workflowDefinitionVersionId = trigger.WorkflowDefinitionVersionId; + return await workflowDefinitionService.FindWorkflowAsync(workflowDefinitionVersionId, cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Http/Services/HttpWorkflowsCacheManager.cs b/src/modules/Elsa.Http/Services/HttpWorkflowsCacheManager.cs index fa5954909..7d6a763cd 100644 --- a/src/modules/Elsa.Http/Services/HttpWorkflowsCacheManager.cs +++ b/src/modules/Elsa.Http/Services/HttpWorkflowsCacheManager.cs @@ -1,89 +1,31 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Options; -using Elsa.Http.Bookmarks; +using Elsa.Caching; using Elsa.Http.Contracts; -using Elsa.Workflows.Activities; -using Elsa.Workflows.Contracts; -using Elsa.Workflows.Helpers; -using Elsa.Workflows.Management.Contracts; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Entities; -using Elsa.Workflows.Runtime.Filters; -using Microsoft.Extensions.Caching.Memory; -using Microsoft.Extensions.Options; -using Open.Linq.AsyncExtensions; namespace Elsa.Http.Services; /// -public class HttpWorkflowsCacheManager( - IMemoryCache memoryCache, - ITriggerStore triggerStore, - IWorkflowDefinitionService workflowDefinitionService, - IChangeTokenSignaler changeTokenSignaler, - IOptions cachingOptions) : IHttpWorkflowsCacheManager +public class HttpWorkflowsCacheManager(ICacheManager cache) : IHttpWorkflowsCacheManager { /// - public async Task<(Workflow? Workflow, ICollection Triggers)?> FindWorkflowAsync(string bookmarkHash, CancellationToken cancellationToken = default) - { - var key = $"workflow:{bookmarkHash}"; - return await memoryCache.GetOrCreateAsync(key, async entry => - { - entry.SetSlidingExpiration(cachingOptions.Value.CacheDuration); - entry.AddExpirationToken(changeTokenSignaler.GetToken(GetTriggerChangeTokenKey(bookmarkHash))); - - var triggers = await FindTriggersAsync(bookmarkHash, cancellationToken).ToList(); - - if (triggers.Count > 1) - return (default, triggers); - - var trigger = triggers.SingleOrDefault(); - - if (trigger == null) - return default; - - var workflow = await FindWorkflowAsync(trigger, cancellationToken); - - if (workflow == null) - return default; - - var changeTokenKey = GetWorkflowChangeTokenKey(workflow.Identity.DefinitionId); - var changeToken = changeTokenSignaler.GetToken(changeTokenKey); - entry.AddExpirationToken(changeToken); - - return (workflow, triggers); - }); - } + public ICacheManager Cache => cache; /// public async Task EvictWorkflowAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) { var changeTokenKey = GetWorkflowChangeTokenKey(workflowDefinitionId); - await changeTokenSignaler.TriggerTokenAsync(changeTokenKey, cancellationToken); + await cache.TriggerTokenAsync(changeTokenKey, cancellationToken); } /// public async Task EvictTriggerAsync(string bookmarkHash, CancellationToken cancellationToken = default) { var changeTokenKey = GetTriggerChangeTokenKey(bookmarkHash); - await changeTokenSignaler.TriggerTokenAsync(changeTokenKey, cancellationToken); + await cache.TriggerTokenAsync(changeTokenKey, cancellationToken); } - private async Task> FindTriggersAsync(string bookmarkHash, CancellationToken cancellationToken) - { - var triggerFilter = new TriggerFilter - { - Hash = bookmarkHash - }; - return await triggerStore.FindManyAsync(triggerFilter, cancellationToken); - } + /// + public string GetWorkflowChangeTokenKey(string workflowDefinitionId) => $"{GetType().FullName}:workflow:{workflowDefinitionId}:changeToken"; - private async Task FindWorkflowAsync(StoredTrigger trigger, CancellationToken cancellationToken) - { - var workflowDefinitionVersionId = trigger.WorkflowDefinitionVersionId; - return await workflowDefinitionService.FindWorkflowAsync(workflowDefinitionVersionId, cancellationToken); - } - - private string GetWorkflowChangeTokenKey(string workflowDefinitionId) => $"{GetType().FullName}:workflow:{workflowDefinitionId}:changeToken"; - private string GetTriggerChangeTokenKey(string bookmarkHash) => $"{GetType().FullName}:trigger:{bookmarkHash}:changeToken"; + /// + public string GetTriggerChangeTokenKey(string bookmarkHash) => $"{GetType().FullName}:trigger:{bookmarkHash}:changeToken"; } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs index ddb0cf7ae..a4af1a8cd 100644 --- a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs +++ b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs @@ -129,8 +129,8 @@ public class WorkflowsFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() - .AddScoped() - .AddScoped() + .AddSingleton() + .AddSingleton() .AddSingleton(IdentityGenerator) .AddScoped(sp => ActivatorUtilities.CreateInstance(sp)) .AddSingleton() diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs index 4ccce4073..68601ea5d 100644 --- a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs @@ -1,5 +1,3 @@ -using System.Diagnostics.CodeAnalysis; -using Elsa.Common.Contracts; using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Activities; @@ -8,39 +6,35 @@ using Elsa.Workflows.Models; using Elsa.Workflows.Notifications; using Elsa.Workflows.Options; using Elsa.Workflows.State; -using Microsoft.Extensions.DependencyInjection; namespace Elsa.Workflows.Services; /// public class WorkflowRunner : IWorkflowRunner { - private readonly IServiceScopeFactory _serviceScopeFactory; + private readonly IServiceProvider _serviceProvider; private readonly IWorkflowExecutionPipeline _pipeline; private readonly IWorkflowStateExtractor _workflowStateExtractor; private readonly IWorkflowBuilderFactory _workflowBuilderFactory; private readonly IIdentityGenerator _identityGenerator; - private readonly ISystemClock _systemClock; private readonly INotificationSender _notificationSender; /// /// Constructor. /// public WorkflowRunner( - IServiceScopeFactory serviceScopeFactory, + IServiceProvider serviceProvider, IWorkflowExecutionPipeline pipeline, IWorkflowStateExtractor workflowStateExtractor, IWorkflowBuilderFactory workflowBuilderFactory, IIdentityGenerator identityGenerator, - ISystemClock systemClock, INotificationSender notificationSender) { - _serviceScopeFactory = serviceScopeFactory; + _serviceProvider = serviceProvider; _pipeline = pipeline; _workflowStateExtractor = workflowStateExtractor; _workflowBuilderFactory = workflowBuilderFactory; _identityGenerator = identityGenerator; - _systemClock = systemClock; _notificationSender = notificationSender; } @@ -90,9 +84,6 @@ public class WorkflowRunner : IWorkflowRunner /// public async Task RunAsync(Workflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) { - // Create a child scope. - using var scope = _serviceScopeFactory.CreateScope(); - // Set up a workflow execution context. var instanceId = options?.WorkflowInstanceId ?? _identityGenerator.GenerateId(); var input = options?.Input; @@ -102,7 +93,7 @@ public class WorkflowRunner : IWorkflowRunner var parentWorkflowInstanceId = options?.ParentWorkflowInstanceId; var statusUpdatedCallback = options?.StatusUpdatedCallback; var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync( - scope.ServiceProvider, + _serviceProvider, workflow, instanceId, correlationId, @@ -123,9 +114,6 @@ public class WorkflowRunner : IWorkflowRunner /// public async Task RunAsync(Workflow workflow, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) { - // Create a child scope. - using var scope = _serviceScopeFactory.CreateScope(); - // Create workflow execution context. var input = options?.Input; var properties = options?.Properties; @@ -134,7 +122,7 @@ public class WorkflowRunner : IWorkflowRunner var parentWorkflowInstanceId = options?.ParentWorkflowInstanceId; var statusUpdatedCallback = options?.StatusUpdatedCallback; var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync( - scope.ServiceProvider, + _serviceProvider, workflow, workflowState, correlationId, diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs index 88420d55a..7e9d1c427 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs @@ -139,9 +139,9 @@ public class WorkflowDefinitionActivity : Composite, IInitializable } } - private async Task FindWorkflowDefinitionAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) + private async Task FindWorkflowAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) { - var workflowDefinitionStore = serviceProvider.GetRequiredService(); + var workflowDefinitionService = serviceProvider.GetRequiredService(); var filter = new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId }; if (!string.IsNullOrWhiteSpace(WorkflowDefinitionVersionId)) @@ -149,32 +149,28 @@ public class WorkflowDefinitionActivity : Composite, IInitializable else filter.VersionOptions = VersionOptions.SpecificVersion(Version); - var workflowDefinition = - await workflowDefinitionStore.FindAsync(filter, cancellationToken) - ?? (await workflowDefinitionStore.FindAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Published }, cancellationToken) - ?? await workflowDefinitionStore.FindAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Latest }, cancellationToken)); + var workflow = + await workflowDefinitionService.FindWorkflowAsync(filter, cancellationToken) + ?? (await workflowDefinitionService.FindWorkflowAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Published }, cancellationToken) + ?? await workflowDefinitionService.FindWorkflowAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Latest }, cancellationToken)); - return workflowDefinition; + return workflow; } async ValueTask IInitializable.InitializeAsync(InitializationContext context) { var serviceProvider = context.ServiceProvider; var cancellationToken = context.CancellationToken; - var workflowDefinition = await FindWorkflowDefinitionAsync(serviceProvider, cancellationToken); + var workflow = await FindWorkflowAsync(serviceProvider, cancellationToken); - if (workflowDefinition == null) + if (workflow == null) throw new Exception($"Could not find workflow definition with ID {WorkflowDefinitionId}."); - // Construct the root activity stored in the activity definitions. - var materializer = serviceProvider.GetRequiredService(); - var root = await materializer.MaterializeAsync(workflowDefinition, cancellationToken); - // Declare input and output variables. - DeclareInputAsVariables(serviceProvider, (inputDescriptor, variable) => Variables.Declare(variable)); - DeclareOutputAsVariables(serviceProvider, (outputDescriptor, variable) => Variables.Declare(variable)); + DeclareInputAsVariables(serviceProvider, (_, variable) => Variables.Declare(variable)); + DeclareOutputAsVariables(serviceProvider, (_, variable) => Variables.Declare(variable)); // Set the root activity. - Root = root; + Root = workflow; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionCacheManager.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionCacheManager.cs index ef87b5c02..c74b3dd13 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionCacheManager.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionCacheManager.cs @@ -1,12 +1,23 @@ +using Elsa.Caching; using Elsa.Common.Models; +using Elsa.Workflows.Management.Filters; namespace Elsa.Workflows.Management.Contracts; +/// +/// Specifies the contract for managing the cache of workflow definitions. +/// public interface IWorkflowDefinitionCacheManager { + /// + /// Gets the cache manager. + /// + ICacheManager Cache { get; } string CreateWorkflowDefinitionVersionCacheKey(string definitionId, VersionOptions versionOptions); + string CreateWorkflowDefinitionFilterCacheKey(WorkflowDefinitionFilter filter); string CreateWorkflowVersionCacheKey(string definitionId, VersionOptions versionOptions); string CreateWorkflowVersionCacheKey(string definitionVersionId); + string CreateWorkflowFilterCacheKey(WorkflowDefinitionFilter filter); string CreateWorkflowDefinitionVersionCacheKey(string definitionVersionId); string CreateWorkflowDefinitionChangeTokenKey(string definitionId); Task EvictWorkflowDefinitionAsync(string definitionId, CancellationToken cancellationToken = default); diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs index 45058401b..da5b416af 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs @@ -1,6 +1,7 @@ using Elsa.Common.Models; using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; namespace Elsa.Workflows.Management.Contracts; @@ -24,6 +25,11 @@ public interface IWorkflowDefinitionService /// Task FindWorkflowDefinitionAsync(string definitionVersionId, CancellationToken cancellationToken = default); + /// + /// Looks for a by the specified . + /// + Task FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); + /// /// Looks for a by the specified definition ID and . /// @@ -33,4 +39,9 @@ public interface IWorkflowDefinitionService /// Looks for a by the specified version ID. /// Task FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default); + + /// + /// Looks for a by the specified . + /// + Task FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Elsa.Workflows.Management.csproj b/src/modules/Elsa.Workflows.Management/Elsa.Workflows.Management.csproj index b7e16da97..99b09454a 100644 --- a/src/modules/Elsa.Workflows.Management/Elsa.Workflows.Management.csproj +++ b/src/modules/Elsa.Workflows.Management/Elsa.Workflows.Management.csproj @@ -10,11 +10,12 @@ + all runtime; build; native; contentfiles; analyzers; buildtransitive - + diff --git a/src/modules/Elsa.Workflows.Management/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Workflows.Management/Extensions/ModuleExtensions.cs index fb9bde99e..d3d177e19 100644 --- a/src/modules/Elsa.Workflows.Management/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.Workflows.Management/Extensions/ModuleExtensions.cs @@ -84,7 +84,7 @@ public static class ModuleExtensions /// /// Adds caching stores feature to the workflow management feature. /// - public static WorkflowManagementFeature UseCachingStores(this WorkflowManagementFeature feature, Action? configure = default) + public static WorkflowManagementFeature UseCache(this WorkflowManagementFeature feature, Action? configure = default) { feature.Module.Configure(configure); return feature; diff --git a/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs b/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs index 889c3ca23..dcf65e950 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/EvictWorkflowDefinitionServiceCache.cs @@ -12,7 +12,7 @@ namespace Elsa.Workflows.Management.Handlers; /// This service listens for specific notifications and triggers cache invalidation operations accordingly. /// [UsedImplicitly] -public class EvictWorkflowDefinitionServiceCache(IWorkflowDefinitionCacheManager workflowDefinitionCacheManager) : +internal class EvictWorkflowDefinitionServiceCache(IWorkflowDefinitionCacheManager workflowDefinitionCacheManager) : INotificationHandler, INotificationHandler, INotificationHandler, diff --git a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs index 873d58f39..9022683be 100644 --- a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs @@ -1,12 +1,10 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Options; using Elsa.Common.Models; using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; using JetBrains.Annotations; using Microsoft.Extensions.Caching.Memory; -using Microsoft.Extensions.Options; namespace Elsa.Workflows.Management.Services; @@ -14,12 +12,7 @@ namespace Elsa.Workflows.Management.Services; /// Decorates an with caching capabilities. /// [UsedImplicitly] -public class CachingWorkflowDefinitionService( - IWorkflowDefinitionService decoratedService, - IWorkflowDefinitionCacheManager cacheManager, - IMemoryCache memoryCache, - IChangeTokenSignaler changeTokenSignaler, - IOptions cachingOptions) : IWorkflowDefinitionService +public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager) : IWorkflowDefinitionService { /// public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) @@ -46,6 +39,15 @@ public class CachingWorkflowDefinitionService( x => x.DefinitionId); } + /// + public async Task FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var cacheKey = cacheManager.CreateWorkflowDefinitionFilterCacheKey(filter); + return await GetFromCacheAsync(cacheKey, + () => decoratedService.FindWorkflowDefinitionAsync(filter, cancellationToken), + x => x.DefinitionId); + } + /// public async Task FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { @@ -64,11 +66,21 @@ public class CachingWorkflowDefinitionService( x => x.Identity.DefinitionId); } + /// + public async Task FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var cacheKey = cacheManager.CreateWorkflowFilterCacheKey(filter); + return await GetFromCacheAsync(cacheKey, + () => decoratedService.FindWorkflowAsync(filter, cancellationToken), + x => x.Identity.DefinitionId); + } + private async Task GetFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc) { - return await memoryCache.GetOrCreateAsync(cacheKey, async entry => + var cache = cacheManager.Cache; + return await cache.GetOrCreateAsync(cacheKey, async entry => { - entry.SetAbsoluteExpiration(cachingOptions.Value.CacheDuration); + entry.SetAbsoluteExpiration(cache.CachingOptions.Value.CacheDuration); var obj = await getObjectFunc(); if (obj == null) @@ -76,7 +88,7 @@ public class CachingWorkflowDefinitionService( var changeTokenKeyInput = getChangeTokenKeyFunc(obj); var changeTokenKey = cacheManager.CreateWorkflowDefinitionChangeTokenKey(changeTokenKeyInput); - entry.AddExpirationToken(changeTokenSignaler.GetToken(changeTokenKey)); + entry.AddExpirationToken(cache.GetToken(changeTokenKey)); return obj; }); } diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionCacheManager.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionCacheManager.cs index 1630c4a6f..8c31d555d 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionCacheManager.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionCacheManager.cs @@ -1,21 +1,40 @@ -using Elsa.Caching.Contracts; +using Elsa.Caching; using Elsa.Common.Models; +using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Filters; namespace Elsa.Workflows.Management.Services; /// -public class WorkflowDefinitionCacheManager(IChangeTokenSignaler changeTokenSignaler) : IWorkflowDefinitionCacheManager +public class WorkflowDefinitionCacheManager(ICacheManager cache, IHasher hasher) : IWorkflowDefinitionCacheManager { + /// + public ICacheManager Cache => cache; + /// public string CreateWorkflowDefinitionVersionCacheKey(string definitionId, VersionOptions versionOptions) => $"WorkflowDefinition:{definitionId}:{versionOptions}"; + /// + public string CreateWorkflowDefinitionFilterCacheKey(WorkflowDefinitionFilter filter) + { + var hash = hasher.Hash(filter); + return $"WorkflowDefinition:{hash}"; + } + /// public string CreateWorkflowVersionCacheKey(string definitionId, VersionOptions versionOptions) => $"Workflow:{definitionId}:{versionOptions}"; /// public string CreateWorkflowVersionCacheKey(string definitionVersionId) => $"Workflow:{definitionVersionId}"; + /// + public string CreateWorkflowFilterCacheKey(WorkflowDefinitionFilter filter) + { + var hash = hasher.Hash(filter); + return $"Workflow:{hash}"; + } + /// public string CreateWorkflowDefinitionVersionCacheKey(string definitionVersionId) => $"WorkflowDefinition:{definitionVersionId}"; @@ -26,6 +45,6 @@ public class WorkflowDefinitionCacheManager(IChangeTokenSignaler changeTokenSign public async Task EvictWorkflowDefinitionAsync(string definitionId, CancellationToken cancellationToken = default) { var changeTokenKey = CreateWorkflowDefinitionChangeTokenKey(definitionId); - await changeTokenSignaler.TriggerTokenAsync(changeTokenKey, cancellationToken); + await cache.TriggerTokenAsync(changeTokenKey, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs index 7065ce80e..8a0b625ff 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs @@ -69,6 +69,12 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService return await _workflowDefinitionStore.FindAsync(filter, cancellationToken); } + /// + public async Task FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + return await _workflowDefinitionStore.FindAsync(filter, cancellationToken); + } + /// public async Task FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { @@ -90,4 +96,15 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService return await MaterializeWorkflowAsync(definition, cancellationToken); } + + /// + public async Task FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken); + + if (definition == null) + return null; + + return await MaterializeWorkflowAsync(definition, cancellationToken); + } } \ 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 68ad723a4..b9e1a83e1 100644 --- a/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs +++ b/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs @@ -1,5 +1,4 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Options; +using Elsa.Caching; using Elsa.Common.Models; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; @@ -7,19 +6,13 @@ using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Models; using Microsoft.Extensions.Caching.Memory; -using Microsoft.Extensions.Options; namespace Elsa.Workflows.Management.Stores; /// /// A decorator for that caches workflow definitions. /// -public class CachingWorkflowDefinitionStore( - IWorkflowDefinitionStore decoratedStore, - IMemoryCache cache, - IHasher hasher, - IChangeTokenSignaler changeTokenSignaler, - IOptions cachingOptions) : IWorkflowDefinitionStore +public class CachingWorkflowDefinitionStore(IWorkflowDefinitionStore decoratedStore, ICacheManager cacheManager, IHasher hasher) : IWorkflowDefinitionStore { private static readonly string CacheInvalidationTokenKey = typeof(CachingWorkflowDefinitionStore).FullName!; @@ -55,14 +48,14 @@ public class CachingWorkflowDefinitionStore( public async Task> FindManyAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var cacheKey = hasher.Hash(filter); - return (await GetOrCreateAsync(cacheKey, () => decoratedStore.FindManyAsync(filter, cancellationToken)))!; + return (await GetOrCreateAsync(cacheKey, async () => (await decoratedStore.FindManyAsync(filter, cancellationToken)).ToList()))!; } /// public async Task> FindManyAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, CancellationToken cancellationToken = default) { var cacheKey = hasher.Hash(filter, order); - return (await GetOrCreateAsync(cacheKey, () => decoratedStore.FindManyAsync(filter, order, cancellationToken)))!; + return (await GetOrCreateAsync(cacheKey, async () => (await decoratedStore.FindManyAsync(filter, order, cancellationToken)).ToList()))!; } /// @@ -104,21 +97,21 @@ public class CachingWorkflowDefinitionStore( public async Task SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { await decoratedStore.SaveAsync(definition, cancellationToken); - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); } /// public async Task SaveManyAsync(IEnumerable definitions, CancellationToken cancellationToken = default) { await decoratedStore.SaveManyAsync(definitions, cancellationToken); - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); } /// public async Task DeleteAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var result = await decoratedStore.DeleteAsync(filter, cancellationToken); - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); return result; } @@ -146,11 +139,11 @@ public class CachingWorkflowDefinitionStore( private async Task GetOrCreateAsync(string key, Func> factory) { var internalKey = $"{typeof(T).Name}:{key}"; - return await cache.GetOrCreateAsync(internalKey, async entry => + return await cacheManager.GetOrCreateAsync(internalKey, async entry => { - var invalidationRequestToken = changeTokenSignaler.GetToken(CacheInvalidationTokenKey); + var invalidationRequestToken = cacheManager.GetToken(CacheInvalidationTokenKey); entry.AddExpirationToken(invalidationRequestToken); - entry.SetAbsoluteExpiration(cachingOptions.Value.CacheDuration); + entry.SetAbsoluteExpiration(cacheManager.CachingOptions.Value.CacheDuration); return await factory(); }); } diff --git a/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerEqualityComparer.cs b/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerEqualityComparer.cs index 4959d120b..100f91daf 100644 --- a/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerEqualityComparer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerEqualityComparer.cs @@ -50,6 +50,7 @@ public class WorkflowTriggerEqualityComparer : IEqualityComparer storedTrigger.Name, storedTrigger.ActivityId, storedTrigger.WorkflowDefinitionId, + storedTrigger.WorkflowDefinitionVersionId, storedTrigger.Hash }; return JsonSerializer.Serialize(input, _settings); diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/ModuleExtensions.cs index f3b3721ba..2a836fe35 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ModuleExtensions.cs @@ -67,7 +67,7 @@ public static class ModuleExtensions /// /// Adds caching stores feature to the workflow runtime feature. /// - public static WorkflowRuntimeFeature UseCachingStores(this WorkflowRuntimeFeature feature, Action? configure = default) + public static WorkflowRuntimeFeature UseCache(this WorkflowRuntimeFeature feature, Action? configure = default) { feature.Module.Configure(configure); return feature; diff --git a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs index e75353436..ed07f9ac8 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs @@ -143,50 +143,24 @@ public class TriggerIndexer : ITriggerIndexer var context = new WorkflowIndexingContext(workflow, cancellationToken); var nodes = await _activityVisitor.VisitAsync(workflow.Root, cancellationToken); - // Get a list of activities that are configured as "startable". - var startableNodes = nodes + // Get a list of trigger activities that are configured as "startable". + var triggerActivities = nodes .Flatten() - .Where(x => x.Activity.GetCanStartWorkflow()) + .Where(x => x.Activity.GetCanStartWorkflow() && x.Activity is ITrigger) + .Select(x => x.Activity) + .Cast() .ToList(); - // For each startable node, create triggers. - foreach (var node in startableNodes) + // For each trigger activity, create a trigger. + foreach (var triggerActivity in triggerActivities) { - var triggers = await GetTriggersAsync(context, node.Activity); + var triggers = await CreateWorkflowTriggersAsync(context, triggerActivity); foreach (var trigger in triggers) yield return trigger; } } - private async Task> GetTriggersAsync(WorkflowIndexingContext context, IActivity activity) - { - // If the activity implements ITrigger, request its trigger data. Otherwise. - if (activity is ITrigger trigger) - return await CreateWorkflowTriggersAsync(context, trigger); - - // Else, create a single workflow trigger with no additional data. - var simpleTrigger = CreateWorkflowTrigger(context, activity); - - return new[] - { - simpleTrigger - }; - } - - private StoredTrigger CreateWorkflowTrigger(WorkflowIndexingContext context, IActivity activity) - { - var workflow = context.Workflow; - return new StoredTrigger - { - Id = _identityGenerator.GenerateId(), - WorkflowDefinitionId = workflow.Identity.DefinitionId, - WorkflowDefinitionVersionId = workflow.Identity.Id, - Name = activity.Type, - ActivityId = activity.Id - }; - } - private async Task> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger trigger) { var workflow = context.Workflow; diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/CachingBookmarkStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/CachingBookmarkStore.cs index 1a8bc66e4..c28e2c609 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/CachingBookmarkStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/CachingBookmarkStore.cs @@ -1,37 +1,30 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Options; +using Elsa.Caching; using Elsa.Workflows.Contracts; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; using Microsoft.Extensions.Caching.Memory; -using Microsoft.Extensions.Options; namespace Elsa.Workflows.Runtime.Stores; /// /// A decorator for that caches bookmark records. /// -public class CachingBookmarkStore( - IBookmarkStore decoratedStore, - IMemoryCache cache, - IHasher hasher, - IChangeTokenSignaler changeTokenSignaler, - IOptions cachingOptions) : IBookmarkStore +public class CachingBookmarkStore(IBookmarkStore decoratedStore, ICacheManager cacheManager, IHasher hasher) : IBookmarkStore { private static readonly string CacheInvalidationTokenKey = typeof(CachingBookmarkStore).FullName!; /// public async ValueTask SaveAsync(StoredBookmark record, CancellationToken cancellationToken = default) { - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); await decoratedStore.SaveAsync(record, cancellationToken); } /// public async ValueTask SaveManyAsync(IEnumerable records, CancellationToken cancellationToken) { - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); await decoratedStore.SaveManyAsync(records, cancellationToken); } @@ -52,17 +45,17 @@ public class CachingBookmarkStore( /// public async ValueTask DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default) { - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); return await decoratedStore.DeleteAsync(filter, cancellationToken); } private async ValueTask GetOrCreateAsync(string key, Func> factory) { - var cacheEntry = await cache.GetOrCreateAsync(key, async entry => + var cacheEntry = await cacheManager.GetOrCreateAsync(key, async entry => { - var invalidationRequestToken = changeTokenSignaler.GetToken(CacheInvalidationTokenKey); + var invalidationRequestToken = cacheManager.GetToken(CacheInvalidationTokenKey); entry.AddExpirationToken(invalidationRequestToken); - entry.SetAbsoluteExpiration(cachingOptions.Value.CacheDuration); + entry.SetAbsoluteExpiration(cacheManager.CachingOptions.Value.CacheDuration); return await factory(); }); diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs index 2dabffbb4..0ebba4a6e 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/CachingTriggerStore.cs @@ -1,37 +1,30 @@ -using Elsa.Caching.Contracts; -using Elsa.Caching.Options; +using Elsa.Caching; using Elsa.Workflows.Contracts; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; using Microsoft.Extensions.Caching.Memory; -using Microsoft.Extensions.Options; namespace Elsa.Workflows.Runtime.Stores; /// /// A decorator for that caches trigger records. /// -public class CachingTriggerStore( - ITriggerStore decoratedStore, - IMemoryCache cache, - IHasher hasher, - IChangeTokenSignaler changeTokenSignaler, - IOptions cachingOptions) : ITriggerStore +public class CachingTriggerStore(ITriggerStore decoratedStore, ICacheManager cacheManager, IHasher hasher) : ITriggerStore { private static readonly string CacheInvalidationTokenKey = typeof(CachingTriggerStore).FullName!; /// public async ValueTask SaveAsync(StoredTrigger record, CancellationToken cancellationToken = default) { - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); await decoratedStore.SaveAsync(record, cancellationToken); } /// public async ValueTask SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default) { - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); await decoratedStore.SaveManyAsync(records, cancellationToken); } @@ -52,25 +45,25 @@ public class CachingTriggerStore( /// public async ValueTask ReplaceAsync(IEnumerable removed, IEnumerable added, CancellationToken cancellationToken = default) { - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); await decoratedStore.ReplaceAsync(removed, added, cancellationToken); } /// public async ValueTask DeleteManyAsync(TriggerFilter filter, CancellationToken cancellationToken = default) { - await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); return await decoratedStore.DeleteManyAsync(filter, cancellationToken); } private async Task GetOrCreateAsync(string key, Func> factory) { var internalKey = $"{typeof(T).Name}:{key}"; - return await cache.GetOrCreateAsync(internalKey, async entry => + return await cacheManager.GetOrCreateAsync(internalKey, async entry => { - var invalidationRequestToken = changeTokenSignaler.GetToken(CacheInvalidationTokenKey); + var invalidationRequestToken = cacheManager.GetToken(CacheInvalidationTokenKey); entry.AddExpirationToken(invalidationRequestToken); - entry.SetAbsoluteExpiration(cachingOptions.Value.CacheDuration); + entry.SetAbsoluteExpiration(cacheManager.CachingOptions.Value.CacheDuration); return await factory(); }); }