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.
This commit is contained in:
parent
b3ae33c586
commit
21c54f583e
1
.github/workflows/packages.yml
vendored
1
.github/workflows/packages.yml
vendored
|
|
@ -6,6 +6,7 @@ on:
|
|||
- 'main'
|
||||
- 'feature/*'
|
||||
- 'issue/*'
|
||||
- 'bug/*'
|
||||
- 'enhancement/*'
|
||||
- 'patch/*'
|
||||
- 'fix/*'
|
||||
|
|
|
|||
|
|
@ -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<MemoryWorkflowInboxMessageStore>();
|
||||
}
|
||||
|
||||
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 =>
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -1,4 +1,3 @@
|
|||
using Elsa.Caching.Contracts;
|
||||
using Elsa.Caching.Distributed.Contracts;
|
||||
using JetBrains.Annotations;
|
||||
using Microsoft.Extensions.Primitives;
|
||||
|
|
|
|||
|
|
@ -1,4 +1,3 @@
|
|||
using Elsa.Caching.Contracts;
|
||||
using Elsa.Caching.Distributed.Contracts;
|
||||
|
||||
namespace Elsa.Caching.Distributed.Services;
|
||||
|
|
|
|||
32
src/modules/Elsa.Caching/Contracts/ICacheManager.cs
Normal file
32
src/modules/Elsa.Caching/Contracts/ICacheManager.cs
Normal file
|
|
@ -0,0 +1,32 @@
|
|||
using Elsa.Caching.Options;
|
||||
using Microsoft.Extensions.Caching.Memory;
|
||||
using Microsoft.Extensions.Options;
|
||||
using Microsoft.Extensions.Primitives;
|
||||
|
||||
namespace Elsa.Caching;
|
||||
|
||||
/// <summary>
|
||||
/// A thin wrapper around <see cref="IMemoryCache"/>, allowing for centralized handling of cache entries.
|
||||
/// </summary>
|
||||
public interface ICacheManager
|
||||
{
|
||||
/// <summary>
|
||||
/// Provides options for configuring caching.
|
||||
/// </summary>
|
||||
IOptions<CachingOptions> CachingOptions { get; }
|
||||
|
||||
/// <summary>
|
||||
/// Gets a change token for the specified key.
|
||||
/// </summary>
|
||||
IChangeToken GetToken(string key);
|
||||
|
||||
/// <summary>
|
||||
/// Triggers the change token for the specified key.
|
||||
/// </summary>
|
||||
ValueTask TriggerTokenAsync(string key, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Gets an item from the cache, or creates it if it doesn't exist.
|
||||
/// </summary>
|
||||
Task<TItem?> GetOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory);
|
||||
}
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
using Microsoft.Extensions.Primitives;
|
||||
|
||||
namespace Elsa.Caching.Contracts;
|
||||
namespace Elsa.Caching;
|
||||
|
||||
/// <summary>
|
||||
/// Provides change tokens for memory caches, allowing code to evict cache entries by triggering a signal.
|
||||
|
|
|
|||
2
src/modules/Elsa.Caching/Elsa.Caching.csproj.DotSettings
Normal file
2
src/modules/Elsa.Caching/Elsa.Caching.csproj.DotSettings
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
|
||||
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=contracts/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>
|
||||
|
|
@ -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<ICacheManager, CacheManager>();
|
||||
Services.AddSingleton<IChangeTokenSignaler, ChangeTokenSignaler>();
|
||||
}
|
||||
}
|
||||
31
src/modules/Elsa.Caching/Services/CacheManager.cs
Normal file
31
src/modules/Elsa.Caching/Services/CacheManager.cs
Normal file
|
|
@ -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;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class CacheManager(IMemoryCache memoryCache, IChangeTokenSignaler changeTokenSignaler, IOptions<CachingOptions> options) : ICacheManager
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public IOptions<CachingOptions> CachingOptions => options;
|
||||
|
||||
/// <inheritdoc />
|
||||
public IChangeToken GetToken(string key)
|
||||
{
|
||||
return changeTokenSignaler.GetToken(key);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public ValueTask TriggerTokenAsync(string key, CancellationToken cancellationToken = default)
|
||||
{
|
||||
return changeTokenSignaler.TriggerTokenAsync(key, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<TItem?> GetOrCreateAsync<TItem>(object key, Func<ICacheEntry, Task<TItem>> factory)
|
||||
{
|
||||
return await memoryCache.GetOrCreateAsync(key, async entry => await factory(entry));
|
||||
}
|
||||
}
|
||||
|
|
@ -1,5 +1,4 @@
|
|||
using System.Collections.Concurrent;
|
||||
using Elsa.Caching.Contracts;
|
||||
using Microsoft.Extensions.Primitives;
|
||||
|
||||
namespace Elsa.Caching.Services;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,14 @@
|
|||
using Elsa.Http.Models;
|
||||
|
||||
namespace Elsa.Http.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a service that looks up HTTP workflows and triggers.
|
||||
/// </summary>
|
||||
public interface IHttpWorkflowLookupService
|
||||
{
|
||||
/// <summary>
|
||||
/// Finds a workflow and trigger by bookmark hash.
|
||||
/// </summary>
|
||||
Task<HttpWorkflowLookupResult?> FindWorkflowAsync(string bookmarkHash, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -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
|
||||
{
|
||||
/// <summary>
|
||||
/// Finds a cached entry by bookmark hash.
|
||||
/// Gets the cache manager.
|
||||
/// </summary>
|
||||
Task<(Workflow? Workflow, ICollection<StoredTrigger> Triggers)?> FindWorkflowAsync(string bookmarkHash, CancellationToken cancellationToken = default);
|
||||
public ICacheManager Cache { get; }
|
||||
|
||||
/// <summary>
|
||||
/// Evicts a cached entry by its definition ID.
|
||||
/// </summary>
|
||||
Task EvictWorkflowAsync(string workflowDefinitionId, CancellationToken cancellationToken = default);
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// Evicts a cached entry by its bookmark hash.
|
||||
/// </summary>
|
||||
Task EvictTriggerAsync(string bookmarkHash, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Gets the key for a workflow change token.
|
||||
/// </summary>
|
||||
string GetWorkflowChangeTokenKey(string workflowDefinitionId);
|
||||
|
||||
/// <summary>
|
||||
/// Gets the key for a trigger change token.
|
||||
/// </summary>
|
||||
string GetTriggerChangeTokenKey(string bookmarkHash);
|
||||
}
|
||||
|
|
@ -15,7 +15,15 @@ public static class ModuleExtensions
|
|||
public static IModule UseHttp(this IModule module, Action<HttpFeature>? configure = default)
|
||||
{
|
||||
module.Configure(configure);
|
||||
module.Use<HttpJavaScriptFeature>();
|
||||
return module;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Install the <see cref="HttpCacheFeature"/> feature to speed up HTTP workflows. Like, a lot.
|
||||
/// </summary>
|
||||
public static HttpFeature UseCache(this HttpFeature feature, Action<HttpCacheFeature>? configure = default)
|
||||
{
|
||||
feature.Module.Configure(configure);
|
||||
return feature;
|
||||
}
|
||||
}
|
||||
32
src/modules/Elsa.Http/Features/HttpCacheFeature.cs
Normal file
32
src/modules/Elsa.Http/Features/HttpCacheFeature.cs
Normal file
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Installs services related to HTTP caching.
|
||||
/// </summary>
|
||||
[DependsOn(typeof(HttpFeature))]
|
||||
[DependsOn(typeof(CachingWorkflowDefinitionsFeature))]
|
||||
public class HttpCacheFeature : FeatureBase
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public HttpCacheFeature(IModule module) : base(module)
|
||||
{
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override void Apply()
|
||||
{
|
||||
Services
|
||||
.AddSingleton<IHttpWorkflowsCacheManager, HttpWorkflowsCacheManager>()
|
||||
.Decorate<IHttpWorkflowLookupService, CachingHttpWorkflowLookupService>()
|
||||
.AddNotificationHandler<InvalidateHttpWorkflowsCache>();
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
|||
/// <summary>
|
||||
/// Installs services related to HTTP services and activities.
|
||||
/// </summary>
|
||||
[DependsOn(typeof(MemoryCacheFeature))]
|
||||
[DependsOn(typeof(HttpJavaScriptFeature))]
|
||||
public class HttpFeature : FeatureBase
|
||||
{
|
||||
/// <inheritdoc />
|
||||
|
|
@ -151,14 +150,13 @@ public class HttpFeature : FeatureBase
|
|||
.AddScoped<IAbsoluteUrlProvider, DefaultAbsoluteUrlProvider>()
|
||||
.AddScoped<IHttpBookmarkProcessor, HttpBookmarkProcessor>()
|
||||
.AddScoped<IRouteTableUpdater, DefaultRouteTableUpdater>()
|
||||
.AddScoped<IHttpWorkflowsCacheManager, HttpWorkflowsCacheManager>()
|
||||
.AddScoped<IHttpWorkflowLookupService, HttpWorkflowLookupService>()
|
||||
.AddScoped(ContentTypeProvider)
|
||||
.AddHttpContextAccessor()
|
||||
|
||||
// Handlers.
|
||||
.AddRequestHandler<ValidateWorkflowRequestHandler, ValidateWorkflowRequest, ValidateWorkflowResponse>()
|
||||
.AddNotificationHandler<UpdateRouteTable>()
|
||||
.AddNotificationHandler<InvalidateHttpWorkflowsCache>()
|
||||
|
||||
// Content parsers.
|
||||
.AddSingleton<IHttpContentParser, JsonHttpContentParser>()
|
||||
|
|
|
|||
|
|
@ -67,13 +67,13 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions<HttpActivity
|
|||
var cancellationToken = httpContext.RequestAborted;
|
||||
var request = httpContext.Request;
|
||||
var method = request.Method.ToLowerInvariant();
|
||||
var httpEndpointCacheManager = serviceProvider.GetRequiredService<IHttpWorkflowsCacheManager>();
|
||||
var httpWorkflowLookupService = serviceProvider.GetRequiredService<IHttpWorkflowLookupService>();
|
||||
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<HttpActivity
|
|||
var trigger = triggers?.FirstOrDefault();
|
||||
if (trigger != null)
|
||||
{
|
||||
var workflow = cachedWorkflowAndTriggers.Value.Workflow!;
|
||||
var workflow = lookupResult.Workflow!;
|
||||
await StartWorkflowAsync(httpContext, trigger, workflow, input);
|
||||
return;
|
||||
}
|
||||
|
|
@ -189,7 +189,7 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions<HttpActivity
|
|||
var workflowHost = await workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken);
|
||||
if (await AuthorizeAsync(serviceProvider, httpContext, workflowHost.Workflow, bookmarkPayload, cancellationToken))
|
||||
return;
|
||||
|
||||
|
||||
await ExecuteWithinTimeoutAsync(async ct =>
|
||||
{
|
||||
var correlationId = await GetCorrelationIdAsync(serviceProvider, httpContext, ct);
|
||||
|
|
|
|||
6
src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs
Normal file
6
src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs
Normal file
|
|
@ -0,0 +1,6 @@
|
|||
using Elsa.Workflows.Activities;
|
||||
using Elsa.Workflows.Runtime.Entities;
|
||||
|
||||
namespace Elsa.Http.Models;
|
||||
|
||||
public record HttpWorkflowLookupResult(Workflow? Workflow, ICollection<StoredTrigger> Triggers);
|
||||
|
|
@ -0,0 +1,40 @@
|
|||
using Elsa.Http.Contracts;
|
||||
using Elsa.Http.Models;
|
||||
using JetBrains.Annotations;
|
||||
using Microsoft.Extensions.Caching.Memory;
|
||||
|
||||
namespace Elsa.Http.Services;
|
||||
|
||||
/// <summary>
|
||||
/// Represents a caching implementation of the IHttpWorkflowLookupService that retrieves workflows using HTTP workflow lookup service and caches the results.
|
||||
/// </summary>
|
||||
[UsedImplicitly]
|
||||
public class CachingHttpWorkflowLookupService(
|
||||
IHttpWorkflowLookupService decoratedService,
|
||||
IHttpWorkflowsCacheManager cacheManager) : IHttpWorkflowLookupService
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<HttpWorkflowLookupResult?> 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;
|
||||
});
|
||||
}
|
||||
}
|
||||
50
src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs
Normal file
50
src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs
Normal file
|
|
@ -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;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class HttpWorkflowLookupService(ITriggerStore triggerStore, IWorkflowDefinitionService workflowDefinitionService) : IHttpWorkflowLookupService
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<HttpWorkflowLookupResult?> 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<IEnumerable<StoredTrigger>> FindTriggersAsync(string bookmarkHash, CancellationToken cancellationToken)
|
||||
{
|
||||
var triggerFilter = new TriggerFilter
|
||||
{
|
||||
Hash = bookmarkHash
|
||||
};
|
||||
return await triggerStore.FindManyAsync(triggerFilter, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<Workflow?> FindWorkflowAsync(StoredTrigger trigger, CancellationToken cancellationToken)
|
||||
{
|
||||
var workflowDefinitionVersionId = trigger.WorkflowDefinitionVersionId;
|
||||
return await workflowDefinitionService.FindWorkflowAsync(workflowDefinitionVersionId, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class HttpWorkflowsCacheManager(
|
||||
IMemoryCache memoryCache,
|
||||
ITriggerStore triggerStore,
|
||||
IWorkflowDefinitionService workflowDefinitionService,
|
||||
IChangeTokenSignaler changeTokenSignaler,
|
||||
IOptions<CachingOptions> cachingOptions) : IHttpWorkflowsCacheManager
|
||||
public class HttpWorkflowsCacheManager(ICacheManager cache) : IHttpWorkflowsCacheManager
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<(Workflow? Workflow, ICollection<StoredTrigger> 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;
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task EvictWorkflowAsync(string workflowDefinitionId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var changeTokenKey = GetWorkflowChangeTokenKey(workflowDefinitionId);
|
||||
await changeTokenSignaler.TriggerTokenAsync(changeTokenKey, cancellationToken);
|
||||
await cache.TriggerTokenAsync(changeTokenKey, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
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<IEnumerable<StoredTrigger>> FindTriggersAsync(string bookmarkHash, CancellationToken cancellationToken)
|
||||
{
|
||||
var triggerFilter = new TriggerFilter
|
||||
{
|
||||
Hash = bookmarkHash
|
||||
};
|
||||
return await triggerStore.FindManyAsync(triggerFilter, cancellationToken);
|
||||
}
|
||||
/// <inheritdoc />
|
||||
public string GetWorkflowChangeTokenKey(string workflowDefinitionId) => $"{GetType().FullName}:workflow:{workflowDefinitionId}:changeToken";
|
||||
|
||||
private async Task<Workflow?> 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";
|
||||
/// <inheritdoc />
|
||||
public string GetTriggerChangeTokenKey(string bookmarkHash) => $"{GetType().FullName}:trigger:{bookmarkHash}:changeToken";
|
||||
}
|
||||
|
|
@ -129,8 +129,8 @@ public class WorkflowsFeature : FeatureBase
|
|||
.AddScoped<IIdentityGraphService, IdentityGraphService>()
|
||||
.AddScoped<IWorkflowStateExtractor, WorkflowStateExtractor>()
|
||||
.AddScoped<IActivitySchedulerFactory, ActivitySchedulerFactory>()
|
||||
.AddScoped<IHasher, Hasher>()
|
||||
.AddScoped<IBookmarkHasher, BookmarkHasher>()
|
||||
.AddSingleton<IHasher, Hasher>()
|
||||
.AddSingleton<IBookmarkHasher, BookmarkHasher>()
|
||||
.AddSingleton(IdentityGenerator)
|
||||
.AddScoped<IBookmarkPayloadSerializer>(sp => ActivatorUtilities.CreateInstance<BookmarkPayloadSerializer>(sp))
|
||||
.AddSingleton<IActivityDescriber, ActivityDescriber>()
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <inheritdoc />
|
||||
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;
|
||||
|
||||
/// <summary>
|
||||
/// Constructor.
|
||||
/// </summary>
|
||||
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
|
|||
/// <inheritdoc />
|
||||
public async Task<RunWorkflowResult> 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
|
|||
/// <inheritdoc />
|
||||
public async Task<RunWorkflowResult> 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,
|
||||
|
|
|
|||
|
|
@ -139,9 +139,9 @@ public class WorkflowDefinitionActivity : Composite, IInitializable
|
|||
}
|
||||
}
|
||||
|
||||
private async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
|
||||
private async Task<Workflow?> FindWorkflowAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
|
||||
{
|
||||
var workflowDefinitionStore = serviceProvider.GetRequiredService<IWorkflowDefinitionStore>();
|
||||
var workflowDefinitionService = serviceProvider.GetRequiredService<IWorkflowDefinitionService>();
|
||||
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<IWorkflowMaterializer>();
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
|
@ -1,12 +1,23 @@
|
|||
using Elsa.Caching;
|
||||
using Elsa.Common.Models;
|
||||
using Elsa.Workflows.Management.Filters;
|
||||
|
||||
namespace Elsa.Workflows.Management.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Specifies the contract for managing the cache of workflow definitions.
|
||||
/// </summary>
|
||||
public interface IWorkflowDefinitionCacheManager
|
||||
{
|
||||
/// <summary>
|
||||
/// Gets the cache manager.
|
||||
/// </summary>
|
||||
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);
|
||||
|
|
|
|||
|
|
@ -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
|
|||
/// </summary>
|
||||
Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(string definitionVersionId, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Looks for a <see cref="WorkflowDefinition"/> by the specified <see cref="WorkflowDefinitionFilter"/>.
|
||||
/// </summary>
|
||||
Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Looks for a <see cref="Workflow"/> by the specified definition ID and <see cref="VersionOptions"/>.
|
||||
/// </summary>
|
||||
|
|
@ -33,4 +39,9 @@ public interface IWorkflowDefinitionService
|
|||
/// Looks for a <see cref="Workflow"/> by the specified version ID.
|
||||
/// </summary>
|
||||
Task<Workflow?> FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Looks for a <see cref="Workflow"/> by the specified <see cref="WorkflowDefinitionFilter"/>.
|
||||
/// </summary>
|
||||
Task<Workflow?> FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -10,11 +10,12 @@
|
|||
<ItemGroup>
|
||||
<PackageReference Include="Humanizer.Core" />
|
||||
<PackageReference Include="IronCompress" />
|
||||
<PackageReference Include="System.Linq.Dynamic.Core" />
|
||||
<PackageReference Include="PolySharp">
|
||||
<PrivateAssets>all</PrivateAssets>
|
||||
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
|
||||
</PackageReference>
|
||||
<PackageReference Include="System.Linq.Dynamic.Core" />
|
||||
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
|
|
|
|||
|
|
@ -84,7 +84,7 @@ public static class ModuleExtensions
|
|||
/// <summary>
|
||||
/// Adds caching stores feature to the workflow management feature.
|
||||
/// </summary>
|
||||
public static WorkflowManagementFeature UseCachingStores(this WorkflowManagementFeature feature, Action<CachingWorkflowDefinitionsFeature>? configure = default)
|
||||
public static WorkflowManagementFeature UseCache(this WorkflowManagementFeature feature, Action<CachingWorkflowDefinitionsFeature>? configure = default)
|
||||
{
|
||||
feature.Module.Configure(configure);
|
||||
return feature;
|
||||
|
|
|
|||
|
|
@ -12,7 +12,7 @@ namespace Elsa.Workflows.Management.Handlers;
|
|||
/// This service listens for specific notifications and triggers cache invalidation operations accordingly.
|
||||
/// </remarks>
|
||||
[UsedImplicitly]
|
||||
public class EvictWorkflowDefinitionServiceCache(IWorkflowDefinitionCacheManager workflowDefinitionCacheManager) :
|
||||
internal class EvictWorkflowDefinitionServiceCache(IWorkflowDefinitionCacheManager workflowDefinitionCacheManager) :
|
||||
INotificationHandler<WorkflowDefinitionPublished>,
|
||||
INotificationHandler<WorkflowDefinitionRetracted>,
|
||||
INotificationHandler<WorkflowDefinitionDeleted>,
|
||||
|
|
|
|||
|
|
@ -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 <see cref="IWorkflowDefinitionService"/> with caching capabilities.
|
||||
/// </summary>
|
||||
[UsedImplicitly]
|
||||
public class CachingWorkflowDefinitionService(
|
||||
IWorkflowDefinitionService decoratedService,
|
||||
IWorkflowDefinitionCacheManager cacheManager,
|
||||
IMemoryCache memoryCache,
|
||||
IChangeTokenSignaler changeTokenSignaler,
|
||||
IOptions<CachingOptions> cachingOptions) : IWorkflowDefinitionService
|
||||
public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager) : IWorkflowDefinitionService
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<Workflow> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
|
||||
|
|
@ -46,6 +39,15 @@ public class CachingWorkflowDefinitionService(
|
|||
x => x.DefinitionId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var cacheKey = cacheManager.CreateWorkflowDefinitionFilterCacheKey(filter);
|
||||
return await GetFromCacheAsync(cacheKey,
|
||||
() => decoratedService.FindWorkflowDefinitionAsync(filter, cancellationToken),
|
||||
x => x.DefinitionId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<Workflow?> FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
|
||||
{
|
||||
|
|
@ -64,11 +66,21 @@ public class CachingWorkflowDefinitionService(
|
|||
x => x.Identity.DefinitionId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<Workflow?> 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<T?> GetFromCacheAsync<T>(string cacheKey, Func<Task<T?>> getObjectFunc, Func<T, string> 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;
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class WorkflowDefinitionCacheManager(IChangeTokenSignaler changeTokenSignaler) : IWorkflowDefinitionCacheManager
|
||||
public class WorkflowDefinitionCacheManager(ICacheManager cache, IHasher hasher) : IWorkflowDefinitionCacheManager
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public ICacheManager Cache => cache;
|
||||
|
||||
/// <inheritdoc />
|
||||
public string CreateWorkflowDefinitionVersionCacheKey(string definitionId, VersionOptions versionOptions) => $"WorkflowDefinition:{definitionId}:{versionOptions}";
|
||||
|
||||
/// <inheritdoc />
|
||||
public string CreateWorkflowDefinitionFilterCacheKey(WorkflowDefinitionFilter filter)
|
||||
{
|
||||
var hash = hasher.Hash(filter);
|
||||
return $"WorkflowDefinition:{hash}";
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public string CreateWorkflowVersionCacheKey(string definitionId, VersionOptions versionOptions) => $"Workflow:{definitionId}:{versionOptions}";
|
||||
|
||||
/// <inheritdoc />
|
||||
public string CreateWorkflowVersionCacheKey(string definitionVersionId) => $"Workflow:{definitionVersionId}";
|
||||
|
||||
/// <inheritdoc />
|
||||
public string CreateWorkflowFilterCacheKey(WorkflowDefinitionFilter filter)
|
||||
{
|
||||
var hash = hasher.Hash(filter);
|
||||
return $"Workflow:{hash}";
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
@ -69,6 +69,12 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService
|
|||
return await _workflowDefinitionStore.FindAsync(filter, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
return await _workflowDefinitionStore.FindAsync(filter, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<Workflow?> FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
|
||||
{
|
||||
|
|
@ -90,4 +96,15 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService
|
|||
|
||||
return await MaterializeWorkflowAsync(definition, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<Workflow?> FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken);
|
||||
|
||||
if (definition == null)
|
||||
return null;
|
||||
|
||||
return await MaterializeWorkflowAsync(definition, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// A decorator for <see cref="IWorkflowDefinitionStore"/> that caches workflow definitions.
|
||||
/// </summary>
|
||||
public class CachingWorkflowDefinitionStore(
|
||||
IWorkflowDefinitionStore decoratedStore,
|
||||
IMemoryCache cache,
|
||||
IHasher hasher,
|
||||
IChangeTokenSignaler changeTokenSignaler,
|
||||
IOptions<CachingOptions> 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<IEnumerable<WorkflowDefinition>> 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()))!;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<IEnumerable<WorkflowDefinition>> FindManyAsync<TOrderBy>(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder<TOrderBy> 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()))!;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
@ -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);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task SaveManyAsync(IEnumerable<WorkflowDefinition> definitions, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await decoratedStore.SaveManyAsync(definitions, cancellationToken);
|
||||
await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<long> 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<T?> GetOrCreateAsync<T>(string key, Func<Task<T>> 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();
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -50,6 +50,7 @@ public class WorkflowTriggerEqualityComparer : IEqualityComparer<StoredTrigger>
|
|||
storedTrigger.Name,
|
||||
storedTrigger.ActivityId,
|
||||
storedTrigger.WorkflowDefinitionId,
|
||||
storedTrigger.WorkflowDefinitionVersionId,
|
||||
storedTrigger.Hash
|
||||
};
|
||||
return JsonSerializer.Serialize(input, _settings);
|
||||
|
|
|
|||
|
|
@ -67,7 +67,7 @@ public static class ModuleExtensions
|
|||
/// <summary>
|
||||
/// Adds caching stores feature to the workflow runtime feature.
|
||||
/// </summary>
|
||||
public static WorkflowRuntimeFeature UseCachingStores(this WorkflowRuntimeFeature feature, Action<CachingWorkflowRuntimeFeature>? configure = default)
|
||||
public static WorkflowRuntimeFeature UseCache(this WorkflowRuntimeFeature feature, Action<CachingWorkflowRuntimeFeature>? configure = default)
|
||||
{
|
||||
feature.Module.Configure(configure);
|
||||
return feature;
|
||||
|
|
|
|||
|
|
@ -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<ITrigger>()
|
||||
.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<IEnumerable<StoredTrigger>> 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<ICollection<StoredTrigger>> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger trigger)
|
||||
{
|
||||
var workflow = context.Workflow;
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// A decorator for <see cref="IBookmarkStore"/> that caches bookmark records.
|
||||
/// </summary>
|
||||
public class CachingBookmarkStore(
|
||||
IBookmarkStore decoratedStore,
|
||||
IMemoryCache cache,
|
||||
IHasher hasher,
|
||||
IChangeTokenSignaler changeTokenSignaler,
|
||||
IOptions<CachingOptions> cachingOptions) : IBookmarkStore
|
||||
public class CachingBookmarkStore(IBookmarkStore decoratedStore, ICacheManager cacheManager, IHasher hasher) : IBookmarkStore
|
||||
{
|
||||
private static readonly string CacheInvalidationTokenKey = typeof(CachingBookmarkStore).FullName!;
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask SaveAsync(StoredBookmark record, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
await decoratedStore.SaveAsync(record, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask SaveManyAsync(IEnumerable<StoredBookmark> 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(
|
|||
/// <inheritdoc />
|
||||
public async ValueTask<long> 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<T?> GetOrCreateAsync<T>(string key, Func<Task<T?>> 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();
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// A decorator for <see cref="ITriggerStore"/> that caches trigger records.
|
||||
/// </summary>
|
||||
public class CachingTriggerStore(
|
||||
ITriggerStore decoratedStore,
|
||||
IMemoryCache cache,
|
||||
IHasher hasher,
|
||||
IChangeTokenSignaler changeTokenSignaler,
|
||||
IOptions<CachingOptions> cachingOptions) : ITriggerStore
|
||||
public class CachingTriggerStore(ITriggerStore decoratedStore, ICacheManager cacheManager, IHasher hasher) : ITriggerStore
|
||||
{
|
||||
private static readonly string CacheInvalidationTokenKey = typeof(CachingTriggerStore).FullName!;
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask SaveAsync(StoredTrigger record, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
await decoratedStore.SaveAsync(record, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask SaveManyAsync(IEnumerable<StoredTrigger> 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(
|
|||
/// <inheritdoc />
|
||||
public async ValueTask ReplaceAsync(IEnumerable<StoredTrigger> removed, IEnumerable<StoredTrigger> added, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await changeTokenSignaler.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken);
|
||||
await decoratedStore.ReplaceAsync(removed, added, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask<long> 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<T?> GetOrCreateAsync<T>(string key, Func<Task<T>> 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();
|
||||
});
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue