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