Add HTTP cache invalidation on auto-update of consuming workflows
This commit is contained in:
parent
e37a0d4a4b
commit
48523c4ecd
|
|
@ -41,14 +41,14 @@ const bool useDapper = false;
|
|||
const bool useProtoActor = false;
|
||||
const bool useHangfire = false;
|
||||
const bool useQuartz = true;
|
||||
const bool useMassTransit = false;
|
||||
const bool useMassTransit = true;
|
||||
const bool useZipCompression = true;
|
||||
const bool runEFCoreMigrations = true;
|
||||
const bool useMemoryStores = false;
|
||||
const bool useCaching = false;
|
||||
const bool useCaching = true;
|
||||
const bool useReadOnlyMode = false;
|
||||
const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.None;
|
||||
const MassTransitBroker useMassTransitBroker = MassTransitBroker.Memory;
|
||||
const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit;
|
||||
const MassTransitBroker useMassTransitBroker = MassTransitBroker.RabbitMq;
|
||||
|
||||
var builder = WebApplication.CreateBuilder(args);
|
||||
var services = builder.Services;
|
||||
|
|
|
|||
|
|
@ -31,4 +31,9 @@ public interface IHttpWorkflowsCacheManager
|
|||
/// Gets the key for a trigger change token.
|
||||
/// </summary>
|
||||
string GetTriggerChangeTokenKey(string bookmarkHash);
|
||||
|
||||
/// <summary>
|
||||
/// Compute the bookmark hash for a given path and method combination.
|
||||
/// </summary>
|
||||
string ComputeBookmarkHash(string path, string method);
|
||||
}
|
||||
|
|
@ -1,6 +1,10 @@
|
|||
using Elsa.Http.Contracts;
|
||||
using Elsa.Http.Bookmarks;
|
||||
using Elsa.Http.Contracts;
|
||||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Management.Notifications;
|
||||
using Elsa.Workflows.Runtime.Contracts;
|
||||
using Elsa.Workflows.Runtime.Entities;
|
||||
using Elsa.Workflows.Runtime.Filters;
|
||||
using Elsa.Workflows.Runtime.Notifications;
|
||||
using JetBrains.Annotations;
|
||||
|
||||
|
|
@ -10,10 +14,16 @@ namespace Elsa.Http.Handlers;
|
|||
/// A handler that invalidates the HTTP workflows cache when a workflow definition is published, retracted, or deleted or when triggers are indexed.
|
||||
/// </summary>
|
||||
[UsedImplicitly]
|
||||
public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflowsCacheManager) :
|
||||
public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflowsCacheManager,
|
||||
ITriggerStore triggerStore,
|
||||
IHttpWorkflowsCacheManager cacheManager) :
|
||||
INotificationHandler<WorkflowDefinitionPublished>,
|
||||
INotificationHandler<WorkflowDefinitionRetracted>,
|
||||
INotificationHandler<WorkflowDefinitionDeleted>,
|
||||
INotificationHandler<WorkflowDefinitionVersionsUpdated>,
|
||||
INotificationHandler<WorkflowDefinitionDeleted>,
|
||||
INotificationHandler<WorkflowDefinitionsDeleted>,
|
||||
INotificationHandler<WorkflowDefinitionVersionDeleted>,
|
||||
INotificationHandler<WorkflowDefinitionVersionsDeleted>,
|
||||
INotificationHandler<WorkflowTriggersIndexed>
|
||||
{
|
||||
/// <inheritdoc />
|
||||
|
|
@ -28,12 +38,45 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo
|
|||
return InvalidateCacheAsync(notification.WorkflowDefinition.DefinitionId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowDefinitionVersionsUpdated notification, CancellationToken cancellationToken)
|
||||
{
|
||||
foreach (WorkflowDefinitionVersionsUpdate versionDefinition in notification.VersionUpdate)
|
||||
{
|
||||
await InvalidateTriggerCacheForDefinitionVersionAsync(versionDefinition.Id, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task HandleAsync(WorkflowDefinitionDeleted notification, CancellationToken cancellationToken)
|
||||
{
|
||||
return InvalidateCacheAsync(notification.DefinitionId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowDefinitionsDeleted notification, CancellationToken cancellationToken)
|
||||
{
|
||||
foreach (string definitionId in notification.DefinitionIds)
|
||||
{
|
||||
await InvalidateCacheAsync(definitionId);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task HandleAsync(WorkflowDefinitionVersionDeleted notification, CancellationToken cancellationToken)
|
||||
{
|
||||
return InvalidateTriggerCacheForDefinitionVersionAsync(notification.WorkflowDefinition.Id, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowDefinitionVersionsDeleted notification, CancellationToken cancellationToken)
|
||||
{
|
||||
foreach (string id in notification.Ids)
|
||||
{
|
||||
await InvalidateTriggerCacheForDefinitionVersionAsync(id, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowTriggersIndexed notification, CancellationToken cancellationToken)
|
||||
{
|
||||
|
|
@ -51,4 +94,27 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo
|
|||
{
|
||||
await httpWorkflowsCacheManager.EvictWorkflowAsync(workflowDefinitionId);
|
||||
}
|
||||
|
||||
private async Task InvalidateTriggerCacheForDefinitionVersionAsync(string workflowDefinitionVersionId, CancellationToken cancellationToken)
|
||||
{
|
||||
var filter = new TriggerFilter
|
||||
{
|
||||
WorkflowDefinitionVersionId = workflowDefinitionVersionId
|
||||
};
|
||||
var triggers = await triggerStore.FindManyAsync(filter, cancellationToken);
|
||||
|
||||
await InvalidateTriggerCacheAsync(triggers, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task InvalidateTriggerCacheAsync(IEnumerable<StoredTrigger> triggers, CancellationToken cancellationToken)
|
||||
{
|
||||
foreach (StoredTrigger trigger in triggers)
|
||||
{
|
||||
if (trigger?.Payload is HttpEndpointBookmarkPayload httpPayload)
|
||||
{
|
||||
var hash = cacheManager.ComputeBookmarkHash(httpPayload.Path, httpPayload.Method);
|
||||
await httpWorkflowsCacheManager.EvictTriggerAsync(hash, cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,10 +1,13 @@
|
|||
using Elsa.Caching;
|
||||
using Elsa.Http.Bookmarks;
|
||||
using Elsa.Http.Contracts;
|
||||
using Elsa.Workflows.Contracts;
|
||||
using Elsa.Workflows.Helpers;
|
||||
|
||||
namespace Elsa.Http.Services;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class HttpWorkflowsCacheManager(ICacheManager cache) : IHttpWorkflowsCacheManager
|
||||
public class HttpWorkflowsCacheManager(ICacheManager cache, IBookmarkHasher bookmarkHasher) : IHttpWorkflowsCacheManager
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public ICacheManager Cache => cache;
|
||||
|
|
@ -28,4 +31,12 @@ public class HttpWorkflowsCacheManager(ICacheManager cache) : IHttpWorkflowsCach
|
|||
|
||||
/// <inheritdoc />
|
||||
public string GetTriggerChangeTokenKey(string bookmarkHash) => $"{GetType().FullName}:trigger:{bookmarkHash}:changeToken";
|
||||
|
||||
/// <inheritdoc />
|
||||
public string ComputeBookmarkHash(string path, string method)
|
||||
{
|
||||
var bookmarkPayload = new HttpEndpointBookmarkPayload(path, method);
|
||||
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<HttpEndpoint>();
|
||||
return bookmarkHasher.Hash(activityTypeName, bookmarkPayload);
|
||||
}
|
||||
}
|
||||
|
|
@ -47,6 +47,9 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IDistributedWorkf
|
|||
distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionVersionsDeleted(notification.Ids), cancellationToken);
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task HandleAsync(WorkflowDefinitionVersionsUpdated notification, CancellationToken cancellationToken) =>
|
||||
distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionVersionsUpdated(notification.DefinitionsAsActivity), cancellationToken);
|
||||
public async Task HandleAsync(WorkflowDefinitionVersionsUpdated notification, CancellationToken cancellationToken)
|
||||
{
|
||||
var updates = notification.VersionUpdate.ToDictionary(x => x.Id, x => x.UsableAsActivity);
|
||||
await distributedEventsDispatcher.DispatchAsync(new Distributed.WorkflowDefinitionVersionsUpdated(updates), cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -96,6 +96,12 @@ public class ActivityRegistry(IActivityDescriber activityDescriber, IEnumerable<
|
|||
|
||||
private void Add(ActivityDescriptor descriptor, ConcurrentDictionary<(string Type, int Version), ActivityDescriptor> activityDescriptors, ICollection<ActivityDescriptor> providerDescriptors)
|
||||
{
|
||||
if (descriptor is null)
|
||||
{
|
||||
logger.LogError("Unable to add a null descriptor");
|
||||
return;
|
||||
}
|
||||
|
||||
foreach (var modifier in modifiers)
|
||||
modifier.Modify(descriptor);
|
||||
|
||||
|
|
|
|||
|
|
@ -16,7 +16,8 @@ internal class EvictWorkflowDefinitionServiceCache(IWorkflowDefinitionCacheManag
|
|||
INotificationHandler<WorkflowDefinitionPublishing>,
|
||||
INotificationHandler<WorkflowDefinitionRetracting>,
|
||||
INotificationHandler<WorkflowDefinitionDeleting>,
|
||||
INotificationHandler<WorkflowDefinitionsDeleting>
|
||||
INotificationHandler<WorkflowDefinitionsDeleting>,
|
||||
INotificationHandler<WorkflowDefinitionVersionsUpdating>
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowDefinitionPublishing notification, CancellationToken cancellationToken)
|
||||
|
|
@ -42,4 +43,11 @@ internal class EvictWorkflowDefinitionServiceCache(IWorkflowDefinitionCacheManag
|
|||
foreach (var definitionId in notification.DefinitionIds)
|
||||
await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(definitionId, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowDefinitionVersionsUpdating notification, CancellationToken cancellationToken)
|
||||
{
|
||||
foreach (var definition in notification.VersionUpdate)
|
||||
await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(definition.DefinitionId, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -80,9 +80,9 @@ public class RefreshActivityRegistry(IActivityRegistryPopulator activityRegistry
|
|||
/// <inheritdoc />
|
||||
public async Task HandleAsync(WorkflowDefinitionVersionsUpdated notification, CancellationToken cancellationToken)
|
||||
{
|
||||
foreach (KeyValuePair<string, bool> definitionAsActivity in notification.DefinitionsAsActivity)
|
||||
foreach (var definition in notification.VersionUpdate)
|
||||
{
|
||||
await UpdateDefinition(definitionAsActivity.Key, definitionAsActivity.Value);
|
||||
await UpdateDefinition(definition.Id, definition.UsableAsActivity);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,6 @@
|
|||
namespace Elsa.Workflows.Management.Notifications;
|
||||
|
||||
/// <summary>
|
||||
/// Represents an update to a specific workflow definition version.
|
||||
/// </summary>
|
||||
public record WorkflowDefinitionVersionsUpdate(string Id, string DefinitionId, bool UsableAsActivity);
|
||||
|
|
@ -6,6 +6,5 @@ namespace Elsa.Workflows.Management.Notifications;
|
|||
/// <summary>
|
||||
/// A notification that is sent when specific workflow definition versions are updated.
|
||||
/// </summary>
|
||||
/// <param name="DefinitionsAsActivity">A dictionary of Ids combined with if the workflow definition is marked usable as activity.</param>
|
||||
[PublicAPI]
|
||||
public record WorkflowDefinitionVersionsUpdated(IDictionary<string, bool> DefinitionsAsActivity) : INotification;
|
||||
public record WorkflowDefinitionVersionsUpdated(IEnumerable<WorkflowDefinitionVersionsUpdate> VersionUpdate) : INotification;
|
||||
|
|
@ -6,6 +6,5 @@ namespace Elsa.Workflows.Management.Notifications;
|
|||
/// <summary>
|
||||
/// A notification that is sent when specific workflow definition versions are about to be updated.
|
||||
/// </summary>
|
||||
/// <param name="DefinitionsAsActivity">A dictionary of the definition version ID combined with if the workflow is marked as usable as activity.</param>
|
||||
[PublicAPI]
|
||||
public record WorkflowDefinitionVersionsUpdating(IDictionary<string, bool> DefinitionsAsActivity) : INotification;
|
||||
public record WorkflowDefinitionVersionsUpdating(IEnumerable<WorkflowDefinitionVersionsUpdate> VersionUpdate) : INotification;
|
||||
|
|
@ -251,10 +251,11 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
|
|||
|
||||
if (updatedWorkflowDefinitions.Any())
|
||||
{
|
||||
var definitionIdsToUpdate = updatedWorkflowDefinitions.ToDictionary(x => x.Id, x => x.Options.UsableAsActivity.GetValueOrDefault());
|
||||
await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdating(definitionIdsToUpdate), cancellationToken);
|
||||
var definitionVersionsUpdates = updatedWorkflowDefinitions.Select(x =>
|
||||
new WorkflowDefinitionVersionsUpdate(x.Id, x.DefinitionId, x.Options.UsableAsActivity.GetValueOrDefault())).ToList();
|
||||
await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdating(definitionVersionsUpdates), cancellationToken);
|
||||
await _workflowDefinitionStore.SaveManyAsync(updatedWorkflowDefinitions, cancellationToken);
|
||||
await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdated(definitionIdsToUpdate), cancellationToken);
|
||||
await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdated(definitionVersionsUpdates), cancellationToken);
|
||||
}
|
||||
|
||||
return updatedWorkflowDefinitions;
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@
|
|||
<ProjectReference Include="..\..\..\src\bundles\Elsa.Server.Web\Elsa.Server.Web.csproj"/>
|
||||
<ProjectReference Include="..\..\..\src\clients\Elsa.Api.Client\Elsa.Api.Client.csproj"/>
|
||||
<ProjectReference Include="..\..\..\src\common\Elsa.Testing.Shared\Elsa.Testing.Shared.csproj"/>
|
||||
<ProjectReference Include="..\..\..\src\modules\Elsa.Caching.Distributed.MassTransit\Elsa.Caching.Distributed.MassTransit.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
|
|
@ -56,6 +57,15 @@
|
|||
<None Update="Scenarios\CachingAndWorkflowDefinitionActivity\Workflows\workflow-definition-parent.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
<None Update="Scenarios\WorkflowActivities\Workflows\CacheInvalidation_Child.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
<None Update="Scenarios\WorkflowActivities\Workflows\CacheInvalidation_Parent.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
<None Update="Scenarios\WorkflowActivities\Workflows\CacheInvalidation_Childv2.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
|
|||
|
|
@ -0,0 +1,15 @@
|
|||
using Elsa.Caching.Distributed.MassTransit.Messages;
|
||||
using Hangfire.Annotations;
|
||||
using MassTransit;
|
||||
|
||||
namespace Elsa.Workflows.ComponentTests.Consumers;
|
||||
|
||||
[UsedImplicitly]
|
||||
public class TriggerChangeTokenSignalConsumer(ITriggerChangeTokenSignalEvents triggerChangeTokenSignalEvents) : IConsumer<TriggerChangeTokenSignal>
|
||||
{
|
||||
public Task Consume(ConsumeContext<TriggerChangeTokenSignal> context)
|
||||
{
|
||||
triggerChangeTokenSignalEvents.OnChangeTokenSignalTriggered(new TriggerChangeTokenSignalEventArgs(context.Message.Key));
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
}
|
||||
|
|
@ -5,7 +5,7 @@ using MassTransit;
|
|||
namespace Elsa.Workflows.ComponentTests.Consumers;
|
||||
|
||||
[UsedImplicitly]
|
||||
public class WorkflowDefinitionEventHandlers(IWorkflowDefinitionEvents workflowDefinitionEvents) : IConsumer<WorkflowDefinitionDeleted>
|
||||
public class WorkflowDefinitionEventConsumer(IWorkflowDefinitionEvents workflowDefinitionEvents) : IConsumer<WorkflowDefinitionDeleted>
|
||||
{
|
||||
public Task Consume(ConsumeContext<WorkflowDefinitionDeleted> context)
|
||||
{
|
||||
|
|
@ -0,0 +1,8 @@
|
|||
namespace Elsa.Workflows.ComponentTests;
|
||||
|
||||
public interface ITriggerChangeTokenSignalEvents
|
||||
{
|
||||
event EventHandler<TriggerChangeTokenSignalEventArgs> ChangeTokenSignalTriggered;
|
||||
|
||||
void OnChangeTokenSignalTriggered(TriggerChangeTokenSignalEventArgs args);
|
||||
}
|
||||
|
|
@ -0,0 +1,6 @@
|
|||
namespace Elsa.Workflows.ComponentTests;
|
||||
|
||||
public class TriggerChangeTokenSignalEventArgs(string key) : EventArgs
|
||||
{
|
||||
public string Key { get; } = key;
|
||||
}
|
||||
|
|
@ -66,7 +66,8 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
|
|||
elsa.UseMassTransit(massTransit =>
|
||||
{
|
||||
massTransit.UseRabbitMq(rabbitMqConnectionString);
|
||||
massTransit.AddConsumer<WorkflowDefinitionEventHandlers>("elsa-test-workflow-definition-updates", true);
|
||||
massTransit.AddConsumer<WorkflowDefinitionEventConsumer>("elsa-test-workflow-definition-updates", true);
|
||||
massTransit.AddConsumer<TriggerChangeTokenSignalConsumer>("elsa-test-change-token-signal", true);
|
||||
});
|
||||
elsa.UseWorkflowManagement(management =>
|
||||
{
|
||||
|
|
@ -77,6 +78,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
|
|||
elsa.UseWorkflowRuntime(runtime =>
|
||||
{
|
||||
runtime.UseMassTransitDispatcher();
|
||||
runtime.UseCache();
|
||||
});
|
||||
};
|
||||
}
|
||||
|
|
@ -86,6 +88,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
|
|||
services.AddSingleton<ISignalManager, SignalManager>();
|
||||
services.AddSingleton<IWorkflowEvents, WorkflowEvents>();
|
||||
services.AddSingleton<IWorkflowDefinitionEvents, WorkflowDefinitionEvents>();
|
||||
services.AddSingleton<ITriggerChangeTokenSignalEvents, TriggerChangeTokenSignalEvents>();
|
||||
services.AddNotificationHandlersFrom<WorkflowServer>();
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,7 @@
|
|||
namespace Elsa.Workflows.ComponentTests.Services;
|
||||
|
||||
public class TriggerChangeTokenSignalEvents : ITriggerChangeTokenSignalEvents
|
||||
{
|
||||
public event EventHandler<TriggerChangeTokenSignalEventArgs>? ChangeTokenSignalTriggered;
|
||||
public void OnChangeTokenSignalTriggered(TriggerChangeTokenSignalEventArgs args) => ChangeTokenSignalTriggered?.Invoke(this, args);
|
||||
}
|
||||
|
|
@ -0,0 +1,120 @@
|
|||
using Elsa.Common.Models;
|
||||
using Elsa.Http;
|
||||
using Elsa.Http.Bookmarks;
|
||||
using Elsa.Http.Contracts;
|
||||
using Elsa.Workflows.Contracts;
|
||||
using Elsa.Workflows.Helpers;
|
||||
using Elsa.Workflows.Management.Contracts;
|
||||
using Elsa.Workflows.Runtime.Filters;
|
||||
using Microsoft.Extensions.Caching.Memory;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowActivities;
|
||||
|
||||
public class AutoUpdateTests : AppComponentTest
|
||||
{
|
||||
private readonly IMemoryCache _cache;
|
||||
private readonly IBookmarkHasher _bookmarkHasher;
|
||||
private readonly IHasher _hasher;
|
||||
private readonly IWorkflowDefinitionCacheManager _definitionCacheManager;
|
||||
private readonly IWorkflowDefinitionPublisher _publisher;
|
||||
private readonly ISignalManager _signalManager;
|
||||
private readonly ITriggerChangeTokenSignalEvents _changeTokenEvents;
|
||||
private readonly IWorkflowDefinitionManager _definitionManager;
|
||||
private readonly IWorkflowDefinitionService _workflowDefinitionService;
|
||||
|
||||
private readonly IHttpWorkflowsCacheManager _httpCacheManager;
|
||||
private readonly IWorkflowDefinitionCacheManager _workflowCacheManager;
|
||||
|
||||
private string _httpChangeToken;
|
||||
private string _triggerChangeToken;
|
||||
private string _graphChangeToken;
|
||||
|
||||
private static readonly object HttpChangeTokenSignal = new();
|
||||
private static readonly object TriggerChangeTokenSignal = new();
|
||||
private static readonly object GraphChangeTokenSignal = new();
|
||||
|
||||
private const string ParentDefinitionVersionId = "4b584e249fdca951";
|
||||
private const string ParentDefinitionId = "878770f04439a55d";
|
||||
private const string ChildDefinitionId = "f353742a9ef6af4";
|
||||
|
||||
public AutoUpdateTests(App app) : base(app)
|
||||
{
|
||||
_cache = Scope.ServiceProvider.GetRequiredService<IMemoryCache>();
|
||||
_bookmarkHasher = Scope.ServiceProvider.GetRequiredService<IBookmarkHasher>();
|
||||
_hasher = Scope.ServiceProvider.GetRequiredService<IHasher>();
|
||||
_definitionCacheManager = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionCacheManager>();
|
||||
_publisher = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionPublisher>();
|
||||
_definitionManager = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionManager>();
|
||||
_workflowDefinitionService = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionService>();
|
||||
|
||||
_httpCacheManager = Scope.ServiceProvider.GetRequiredService<IHttpWorkflowsCacheManager>();
|
||||
_workflowCacheManager = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionCacheManager>();
|
||||
|
||||
_signalManager = Scope.ServiceProvider.GetRequiredService<ISignalManager>();
|
||||
_changeTokenEvents = Scope.ServiceProvider.GetRequiredService<ITriggerChangeTokenSignalEvents>();
|
||||
_changeTokenEvents.ChangeTokenSignalTriggered += OnChangeTokenSignalTriggered;
|
||||
}
|
||||
|
||||
[Fact(DisplayName = "Updating a workflow with `auto update consuming workflows` should invalidate consuming workflows from cache")]
|
||||
public async Task UpdateWorkflowWithAutoUpdate()
|
||||
{
|
||||
//Run workflow to make sure the all required items for running the workflow are in the cache
|
||||
var client = WorkflowServer.CreateHttpWorkflowClient();
|
||||
await client.GetStringAsync("test-cache-invalidation");
|
||||
|
||||
//Make sure the items are in the cache
|
||||
var hash = ComputeBookmarkHash("/test-cache-invalidation", "get");
|
||||
Assert.True(_cache.TryGetValue($"http-workflow:{hash}", out _));
|
||||
|
||||
var filter = new TriggerFilter
|
||||
{
|
||||
Hash = hash
|
||||
};
|
||||
var hashedFilter = _hasher.Hash(filter);
|
||||
Assert.True(_cache.TryGetValue($"IEnumerable`1:{hashedFilter}", out _));
|
||||
|
||||
var parentVersionCacheKey = _definitionCacheManager.CreateWorkflowVersionCacheKey(ParentDefinitionVersionId);
|
||||
Assert.True(_cache.TryGetValue(parentVersionCacheKey, out _));
|
||||
|
||||
//Set change tokens
|
||||
_httpChangeToken = _workflowCacheManager.CreateWorkflowDefinitionChangeTokenKey(ParentDefinitionId);
|
||||
_triggerChangeToken = _httpCacheManager.GetTriggerChangeTokenKey(hash);
|
||||
_graphChangeToken = _workflowCacheManager.CreateWorkflowDefinitionChangeTokenKey(ParentDefinitionId);
|
||||
|
||||
//(Act) Save the draft version of the child workflow and update the references
|
||||
await _publisher.PublishAsync(ChildDefinitionId);
|
||||
|
||||
//Wait till the notifications for updating the cache have been send and check the cache.
|
||||
await _signalManager.WaitAsync<TriggerChangeTokenSignalEventArgs>(HttpChangeTokenSignal);
|
||||
await _signalManager.WaitAsync<TriggerChangeTokenSignalEventArgs>(TriggerChangeTokenSignal);
|
||||
await _signalManager.WaitAsync<TriggerChangeTokenSignalEventArgs>(GraphChangeTokenSignal);
|
||||
|
||||
Assert.False(_cache.TryGetValue($"http-workflow:{hash}", out _));
|
||||
Assert.False(_cache.TryGetValue($"IEnumerable`1:{hashedFilter}", out _));
|
||||
Assert.False(_cache.TryGetValue(parentVersionCacheKey, out _));
|
||||
}
|
||||
|
||||
private void OnChangeTokenSignalTriggered(object? sender, TriggerChangeTokenSignalEventArgs args)
|
||||
{
|
||||
if (args.Key == _httpChangeToken)
|
||||
{
|
||||
_signalManager.Trigger(HttpChangeTokenSignal, args);
|
||||
}
|
||||
if (args.Key == _triggerChangeToken)
|
||||
{
|
||||
_signalManager.Trigger(TriggerChangeTokenSignal, args);
|
||||
}
|
||||
if (args.Key == _graphChangeToken)
|
||||
{
|
||||
_signalManager.Trigger(GraphChangeTokenSignal, args);
|
||||
}
|
||||
}
|
||||
|
||||
private string ComputeBookmarkHash(string path, string method)
|
||||
{
|
||||
var bookmarkPayload = new HttpEndpointBookmarkPayload(path, method);
|
||||
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<HttpEndpoint>();
|
||||
return _bookmarkHasher.Hash(activityTypeName, bookmarkPayload);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,4 +1,3 @@
|
|||
using Elsa.Workflows.ComponentTests.Scenarios.WorkflowActivities.Workflows;
|
||||
using Elsa.Workflows.Contracts;
|
||||
using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
|
||||
using Elsa.Workflows.Management.Contracts;
|
||||
|
|
@ -18,8 +17,9 @@ public class DeleteWorkflowTests : AppComponentTest
|
|||
public DeleteWorkflowTests(App app) : base(app)
|
||||
{
|
||||
_scope1 = app.Cluster.Pod1.Services.CreateScope();
|
||||
_scope2 = app.Cluster.Pod2.Services.CreateScope();
|
||||
_scope3 = app.Cluster.Pod3.Services.CreateScope();
|
||||
// Disabled these scope creations since this prevents events from firing in other tests.
|
||||
// _scope2 = app.Cluster.Pod2.Services.CreateScope();
|
||||
// _scope3 = app.Cluster.Pod3.Services.CreateScope();
|
||||
_signalManager = Scope.ServiceProvider.GetRequiredService<ISignalManager>();
|
||||
_workflowDefinitionEvents = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionEvents>();
|
||||
_workflowDefinitionEvents.WorkflowDefinitionDeleted += OnWorkflowDefinionDeleted;
|
||||
|
|
@ -36,14 +36,14 @@ public class DeleteWorkflowTests : AppComponentTest
|
|||
WorkflowTypeDeletedFromRegistry(_scope1);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
[Fact(Skip = "Clustered tests are interfering with other event driven tests")]
|
||||
public async Task DeleteWorkflow_Clustered()
|
||||
{
|
||||
EnsureWorkflowInRegistry(_scope1);
|
||||
EnsureWorkflowInRegistry(_scope2);
|
||||
EnsureWorkflowInRegistry(_scope3);
|
||||
|
||||
var workflowDefinitionManager = _scope3.ServiceProvider.GetRequiredService<IWorkflowDefinitionManager>();
|
||||
var workflowDefinitionManager = _scope1.ServiceProvider.GetRequiredService<IWorkflowDefinitionManager>();
|
||||
await workflowDefinitionManager.DeleteByDefinitionIdAsync(Workflows.DeleteWorkflow.DefinitionId);
|
||||
|
||||
WorkflowTypeDeletedFromRegistry(_scope1);
|
||||
|
|
@ -79,8 +79,6 @@ public class DeleteWorkflowTests : AppComponentTest
|
|||
|
||||
protected override void OnDispose()
|
||||
{
|
||||
_scope1.Dispose();
|
||||
_scope2.Dispose();
|
||||
_scope3.Dispose();
|
||||
_workflowDefinitionEvents.WorkflowDefinitionDeleted -= OnWorkflowDefinionDeleted;
|
||||
}
|
||||
}
|
||||
|
|
@ -9,22 +9,20 @@ namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowActivities;
|
|||
|
||||
public class SaveWorkflowTests(App app) : AppComponentTest(app)
|
||||
{
|
||||
private readonly IServiceScope _scope = app.Cluster.Pod1.Services.CreateScope();
|
||||
|
||||
[Theory]
|
||||
[Theory(DisplayName = "Saving workflows updates ActivityRegistry")]
|
||||
[InlineData("Save1", true, true, true, true)]
|
||||
[InlineData("Save2", true, false, true, false)]
|
||||
[InlineData("Save3", false, true, false, false)]
|
||||
[InlineData("Save4", false, false, false, false)]
|
||||
private async Task SaveWorkflow(string name, bool usableAsActivity, bool publish, bool expectedInRegistry, bool isBrowsable)
|
||||
public async Task ActivityRegistry(string name, bool usableAsActivity, bool publish, bool expectedInRegistry, bool isBrowsable)
|
||||
{
|
||||
var activityRegistry = _scope.ServiceProvider.GetRequiredService<IActivityRegistry>();
|
||||
var activityRegistry = Scope.ServiceProvider.GetRequiredService<IActivityRegistry>();
|
||||
|
||||
var descriptor = activityRegistry.Find(name);
|
||||
if (descriptor is not null)
|
||||
activityRegistry.Remove(typeof(WorkflowDefinitionActivityProvider), descriptor);
|
||||
|
||||
var importer = _scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionImporter>();
|
||||
var importer = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionImporter>();
|
||||
var request = new SaveWorkflowDefinitionRequest
|
||||
{
|
||||
Model = new WorkflowDefinitionModel
|
||||
|
|
@ -53,9 +51,4 @@ public class SaveWorkflowTests(App app) : AppComponentTest(app)
|
|||
Assert.Null(descriptor);
|
||||
}
|
||||
}
|
||||
|
||||
protected override void OnDispose()
|
||||
{
|
||||
_scope.Dispose();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,67 @@
|
|||
{
|
||||
"id": "23d18a46dc3000fd",
|
||||
"definitionId": "f353742a9ef6af4",
|
||||
"name": "CacheInvalidation_Child",
|
||||
"createdAt": "2024-05-28T08:10:31.5429337+00:00",
|
||||
"version": 1,
|
||||
"toolVersion": "3.2.0.0",
|
||||
"variables": [],
|
||||
"inputs": [],
|
||||
"outputs": [],
|
||||
"outcomes": [],
|
||||
"customProperties": {},
|
||||
"isReadonly": false,
|
||||
"isSystem": false,
|
||||
"isLatest": true,
|
||||
"isPublished": true,
|
||||
"options": {
|
||||
"usableAsActivity": true,
|
||||
"autoUpdateConsumingWorkflows": true
|
||||
},
|
||||
"root": {
|
||||
"type": "Elsa.Flowchart",
|
||||
"version": 1,
|
||||
"id": "5d146ccb42668b6b",
|
||||
"nodeId": "Workflow1:5d146ccb42668b6b",
|
||||
"metadata": {},
|
||||
"customProperties": {
|
||||
"source": "FlowchartJsonConverter.cs:45",
|
||||
"notFoundConnections": [],
|
||||
"canStartWorkflow": false,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"activities": [
|
||||
{
|
||||
"correlationId": {
|
||||
"typeName": "String",
|
||||
"expression": {
|
||||
"type": "Literal",
|
||||
"value": "Child version 1"
|
||||
}
|
||||
},
|
||||
"id": "d719079757f5d1bf",
|
||||
"nodeId": "Workflow1:5d146ccb42668b6b:d719079757f5d1bf",
|
||||
"name": "Correlate1",
|
||||
"type": "Elsa.Correlate",
|
||||
"version": 1,
|
||||
"customProperties": {
|
||||
"canStartWorkflow": false,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"metadata": {
|
||||
"designer": {
|
||||
"position": {
|
||||
"x": -127.40211486816406,
|
||||
"y": 526.2000122070312
|
||||
},
|
||||
"size": {
|
||||
"width": 132.7375030517578,
|
||||
"height": 49.60000228881836
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
],
|
||||
"connections": []
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,67 @@
|
|||
{
|
||||
"id": "3da9c310a071ac7b",
|
||||
"definitionId": "f353742a9ef6af4",
|
||||
"name": "CacheInvalidation_Child",
|
||||
"createdAt": "2024-05-28T11:00:47.8186457+00:00",
|
||||
"version": 2,
|
||||
"toolVersion": "3.2.0.0",
|
||||
"variables": [],
|
||||
"inputs": [],
|
||||
"outputs": [],
|
||||
"outcomes": [],
|
||||
"customProperties": {},
|
||||
"isReadonly": false,
|
||||
"isSystem": false,
|
||||
"isLatest": true,
|
||||
"isPublished": false,
|
||||
"options": {
|
||||
"usableAsActivity": true,
|
||||
"autoUpdateConsumingWorkflows": true
|
||||
},
|
||||
"root": {
|
||||
"type": "Elsa.Flowchart",
|
||||
"version": 1,
|
||||
"id": "5d146ccb42668b6b",
|
||||
"nodeId": "Workflow1:5d146ccb42668b6b",
|
||||
"metadata": {},
|
||||
"customProperties": {
|
||||
"source": "FlowchartJsonConverter.cs:45",
|
||||
"notFoundConnections": [],
|
||||
"canStartWorkflow": false,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"activities": [
|
||||
{
|
||||
"correlationId": {
|
||||
"typeName": "String",
|
||||
"expression": {
|
||||
"type": "Literal",
|
||||
"value": "Child version 2"
|
||||
}
|
||||
},
|
||||
"id": "d719079757f5d1bf",
|
||||
"nodeId": "Workflow1:5d146ccb42668b6b:d719079757f5d1bf",
|
||||
"name": "Correlate1",
|
||||
"type": "Elsa.Correlate",
|
||||
"version": 1,
|
||||
"customProperties": {
|
||||
"canStartWorkflow": false,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"metadata": {
|
||||
"designer": {
|
||||
"position": {
|
||||
"x": -127.40211486816406,
|
||||
"y": 526.2000122070312
|
||||
},
|
||||
"size": {
|
||||
"width": 132.7375030517578,
|
||||
"height": 49.60000228881836
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
],
|
||||
"connections": []
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,188 @@
|
|||
{
|
||||
"id": "4b584e249fdca951",
|
||||
"definitionId": "878770f04439a55d",
|
||||
"name": "CacheInvalidation_Parent",
|
||||
"createdAt": "2024-05-28T08:11:06.1123531+00:00",
|
||||
"version": 1,
|
||||
"toolVersion": "3.2.0.0",
|
||||
"variables": [
|
||||
{
|
||||
"id": "47ab4de953db266",
|
||||
"name": "CorrId",
|
||||
"typeName": "String",
|
||||
"isArray": false,
|
||||
"storageDriverTypeName": "Elsa.Workflows.Services.WorkflowStorageDriver, Elsa.Workflows.Core"
|
||||
}
|
||||
],
|
||||
"inputs": [],
|
||||
"outputs": [],
|
||||
"outcomes": [],
|
||||
"customProperties": {},
|
||||
"isReadonly": false,
|
||||
"isSystem": false,
|
||||
"isLatest": true,
|
||||
"isPublished": true,
|
||||
"options": {
|
||||
"autoUpdateConsumingWorkflows": false
|
||||
},
|
||||
"root": {
|
||||
"type": "Elsa.Flowchart",
|
||||
"version": 1,
|
||||
"id": "76f03074bbc8cff3",
|
||||
"nodeId": "Workflow2:76f03074bbc8cff3",
|
||||
"metadata": {},
|
||||
"customProperties": {
|
||||
"source": "FlowchartJsonConverter.cs:45",
|
||||
"notFoundConnections": [],
|
||||
"canStartWorkflow": false,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"activities": [
|
||||
{
|
||||
"workflowDefinitionId": "f353742a9ef6af4",
|
||||
"workflowDefinitionVersionId": "23d18a46dc3000fd",
|
||||
"latestAvailablePublishedVersion": 1,
|
||||
"latestAvailablePublishedVersionId": "23d18a46dc3000fd",
|
||||
"id": "b8fe7116016141f",
|
||||
"nodeId": "Workflow2:76f03074bbc8cff3:b8fe7116016141f",
|
||||
"name": "CacheInvalidationChild1",
|
||||
"type": "CacheInvalidationChild",
|
||||
"version": 1,
|
||||
"customProperties": {
|
||||
"canStartWorkflow": false,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"metadata": {
|
||||
"designer": {
|
||||
"position": {
|
||||
"x": -40,
|
||||
"y": 320
|
||||
},
|
||||
"size": {
|
||||
"width": 203.9499969482422,
|
||||
"height": 49.60000228881836
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"value": {
|
||||
"typeName": "String",
|
||||
"expression": {
|
||||
"type": "Literal",
|
||||
"value": "Parent version 1"
|
||||
}
|
||||
},
|
||||
"id": "ad9ee49c807b852e",
|
||||
"nodeId": "Workflow2:76f03074bbc8cff3:ad9ee49c807b852e",
|
||||
"name": "SetName1",
|
||||
"type": "Elsa.SetName",
|
||||
"version": 1,
|
||||
"customProperties": {
|
||||
"canStartWorkflow": false,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"metadata": {
|
||||
"designer": {
|
||||
"position": {
|
||||
"x": -303.3261413574219,
|
||||
"y": 320
|
||||
},
|
||||
"size": {
|
||||
"width": 137.6354217529297,
|
||||
"height": 49.60000228881836
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"path": {
|
||||
"typeName": "String",
|
||||
"expression": {
|
||||
"type": "Literal",
|
||||
"value": "test-cache-invalidation"
|
||||
}
|
||||
},
|
||||
"supportedMethods": {
|
||||
"typeName": "String[]",
|
||||
"expression": {
|
||||
"type": "Object",
|
||||
"value": "[\u0022GET\u0022]"
|
||||
}
|
||||
},
|
||||
"authorize": {
|
||||
"typeName": "Boolean",
|
||||
"expression": {
|
||||
"type": "Literal",
|
||||
"value": false
|
||||
}
|
||||
},
|
||||
"policy": {
|
||||
"typeName": "String",
|
||||
"expression": {
|
||||
"type": "Literal"
|
||||
}
|
||||
},
|
||||
"requestTimeout": null,
|
||||
"requestSizeLimit": null,
|
||||
"fileSizeLimit": null,
|
||||
"allowedFileExtensions": null,
|
||||
"blockedFileExtensions": null,
|
||||
"allowedMimeTypes": null,
|
||||
"exposeRequestTooLargeOutcome": false,
|
||||
"exposeFileTooLargeOutcome": false,
|
||||
"exposeInvalidFileExtensionOutcome": false,
|
||||
"exposeInvalidFileMimeTypeOutcome": false,
|
||||
"parsedContent": null,
|
||||
"files": null,
|
||||
"routeData": null,
|
||||
"queryStringData": null,
|
||||
"headers": null,
|
||||
"result": null,
|
||||
"id": "109b9c25e55375ba",
|
||||
"nodeId": "Workflow2:76f03074bbc8cff3:109b9c25e55375ba",
|
||||
"name": "HttpEndpoint1",
|
||||
"type": "Elsa.HttpEndpoint",
|
||||
"version": 1,
|
||||
"customProperties": {
|
||||
"canStartWorkflow": true,
|
||||
"runAsynchronously": false
|
||||
},
|
||||
"metadata": {
|
||||
"designer": {
|
||||
"position": {
|
||||
"x": -622.4391784667969,
|
||||
"y": 320
|
||||
},
|
||||
"size": {
|
||||
"width": 176.0625,
|
||||
"height": 49.60000228881836
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
],
|
||||
"connections": [
|
||||
{
|
||||
"source": {
|
||||
"activity": "ad9ee49c807b852e",
|
||||
"port": "Done"
|
||||
},
|
||||
"target": {
|
||||
"activity": "b8fe7116016141f",
|
||||
"port": "In"
|
||||
}
|
||||
},
|
||||
{
|
||||
"source": {
|
||||
"activity": "109b9c25e55375ba",
|
||||
"port": "Done"
|
||||
},
|
||||
"target": {
|
||||
"activity": "ad9ee49c807b852e",
|
||||
"port": "In"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
|
|
@ -8,6 +8,7 @@ public class DeleteWorkflow : WorkflowBase
|
|||
{
|
||||
public static readonly string DefinitionId = Guid.NewGuid().ToString();
|
||||
public static readonly string Type = nameof(DeleteWorkflow);
|
||||
|
||||
protected override void Build(IWorkflowBuilder builder)
|
||||
{
|
||||
builder.Name = Type;
|
||||
|
|
|
|||
Loading…
Reference in a new issue