diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs index f5cdf962a..149591861 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs @@ -55,4 +55,9 @@ public interface IWorkflowDefinitionService /// Looks for a by the specified . /// Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); + + /// + /// Looks for all s that match the specified . + /// + Task> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionStore.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionStore.cs index ce4ae3663..911b5ead0 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionStore.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionStore.cs @@ -159,5 +159,5 @@ public interface IWorkflowDefinitionStore /// The name. /// The definition ID to exclude from the check. /// The cancellation token. - Task GetIsNameUnique(string name, string? definitionId = default, CancellationToken cancellationToken = default); + Task GetIsNameUnique(string name, string? definitionId = null, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceQuery.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceQuery.cs new file mode 100644 index 000000000..40ee4d463 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowReferenceQuery.cs @@ -0,0 +1,15 @@ +namespace Elsa.Workflows.Management; + +/// +/// Finds all latest versions of workflow definitions that reference a specific workflow definition. +/// +public interface IWorkflowReferenceQuery +{ + /// + /// Queries all latest versions of workflow definitions that reference the specified workflow definition. + /// + /// The ID of the workflow definition to query references for. + /// The cancellation token to cancel the operation. + /// A collection of workflow definition IDs that reference the specified workflow definition. + Task> ExecuteAsync(string workflowDefinitionId, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index ff369a7c4..30f7c7f18 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -28,6 +28,7 @@ using Elsa.Workflows.Management.Stores; using Elsa.Workflows.Serialization.Serializers; using JetBrains.Annotations; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; namespace Elsa.Workflows.Management.Features; @@ -51,6 +52,7 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module) private const string SystemCategory = "System"; private Func _workflowDefinitionPublisher = sp => ActivatorUtilities.CreateInstance(sp); + private Func _workflowReferenceQuery = sp => ActivatorUtilities.CreateInstance(sp); private string CompressionAlgorithm { get; set; } = nameof(None); private LogPersistenceMode LogPersistenceMode { get; set; } = LogPersistenceMode.Include; @@ -191,11 +193,23 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module) return this; } - public WorkflowManagementFeature WithWorkflowDefinitionPublisher(Func workflowDefinitionPublisher) + public WorkflowManagementFeature UseWorkflowDefinitionPublisher(Func workflowDefinitionPublisher) { _workflowDefinitionPublisher = workflowDefinitionPublisher; return this; } + + public WorkflowManagementFeature UseWorkflowReferenceFinder() where T : class, IWorkflowReferenceQuery + { + Services.TryAddScoped(); + return UseWorkflowReferenceFinder(sp => sp.GetRequiredService()); + } + + public WorkflowManagementFeature UseWorkflowReferenceFinder(Func workflowReferenceFinder) + { + _workflowReferenceQuery = workflowReferenceFinder; + return this; + } /// [RequiresUnreferencedCode("The assembly containing the specified marker type will be scanned for activity types.")] @@ -217,6 +231,7 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module) .AddScoped() .AddScoped() .AddScoped() + .AddScoped(_workflowReferenceQuery) .AddScoped(_workflowDefinitionPublisher) .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs index f8936d8ef..059457ba8 100644 --- a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs @@ -11,7 +11,7 @@ namespace Elsa.Workflows.Management.Services; /// Decorates an with caching capabilities. /// [UsedImplicitly] -public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager) : IWorkflowDefinitionService +public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowDefinitionService { /// public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) @@ -41,10 +41,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat /// public Task FindWorkflowDefinitionAsync(WorkflowDefinitionHandle handle, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter - { - DefinitionHandle = handle - }; + var filter = new WorkflowDefinitionFilter { DefinitionHandle = handle }; return FindWorkflowDefinitionAsync(filter, cancellationToken); } @@ -79,10 +76,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat /// public Task FindWorkflowGraphAsync(WorkflowDefinitionHandle definitionHandle, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter - { - DefinitionHandle = definitionHandle - }; + var filter = new WorkflowDefinitionFilter { DefinitionHandle = definitionHandle }; return FindWorkflowGraphAsync(filter, cancellationToken); } @@ -96,6 +90,23 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat x => x.Workflow.Identity.DefinitionId); } + public async Task> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken); + var workflowGraphs = new List(); + foreach (var workflowDefinition in workflowDefinitions) + { + var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(workflowDefinition.Id); + var workflowGraph = await GetFromCacheAsync( + cacheKey, + async () => await MaterializeWorkflowAsync(workflowDefinition, cancellationToken), + wf => wf.Workflow.Identity.DefinitionId); + workflowGraphs.Add(workflowGraph!); + } + + return workflowGraphs; + } + private async Task GetFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc) { var cache = cacheManager.Cache; diff --git a/src/modules/Elsa.Workflows.Management/Services/DefaultWorkflowReferenceQuery.cs b/src/modules/Elsa.Workflows.Management/Services/DefaultWorkflowReferenceQuery.cs new file mode 100644 index 000000000..b0bf0ffab --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/DefaultWorkflowReferenceQuery.cs @@ -0,0 +1,67 @@ +using Elsa.Common.Models; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.Services; + +public class DefaultWorkflowReferenceQuery(IWorkflowDefinitionService workflowDefinitionService, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowReferenceQuery +{ + public async Task> ExecuteAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) + { + var workflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(workflowDefinitionId, cancellationToken); + return workflowGraphs.Select(x => x.Workflow.Identity.DefinitionId).Distinct(); + } + + private async Task> FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(string updatedWorkflowDefinitionId, CancellationToken cancellationToken) + { + var allWorkflowGraphs = (await GetAllWorkflowGraphsAsync(cancellationToken)).ToList(); + var filteredWorkflowGraphs = allWorkflowGraphs + .Where(x => FindWorkflowActivityDefinitionActivityNodes(x.Root).Any(y => y.WorkflowDefinitionId == updatedWorkflowDefinitionId)) + .ToList(); + + return filteredWorkflowGraphs; + } + + private IEnumerable FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent) + { + foreach (var child in parent.Children) + { + if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity) + yield return workflowDefinitionActivity; + else + { + foreach (var grandChild in FindWorkflowActivityDefinitionActivityNodes(child)) + yield return grandChild; + } + } + } + + private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) + { + var workflowDefinitionFilter = new WorkflowDefinitionFilter + { + VersionOptions = VersionOptions.LatestOrPublished, + IsReadonly = false + }; + var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); + + // If there are workflow definition summaries with the same definition ID, only take the latest version. + workflowDefinitionSummaries = workflowDefinitionSummaries + .GroupBy(x => x.DefinitionId) + .Select(x => x.OrderByDescending(y => y.Version).First()) + .ToList(); + + var workflowGraphs = new List(); + + foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) + { + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.Id, cancellationToken); + + if (workflowGraph != null) + workflowGraphs.Add(workflowGraph); + } + + return workflowGraphs; + } +} \ 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 bd3435846..104ea4e8c 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs @@ -95,4 +95,18 @@ public class WorkflowDefinitionService( return await MaterializeWorkflowAsync(definition, cancellationToken); } + + /// + public async Task> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + { + var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken); + var workflowGraphs = new List(); + foreach (var workflowDefinition in workflowDefinitions) + { + var workflowGraph = await MaterializeWorkflowAsync(workflowDefinition, cancellationToken); + workflowGraphs.Add(workflowGraph); + } + + return workflowGraphs; + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs index 4cdd7a2a4..988fdc7bf 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -1,7 +1,6 @@ using Elsa.Common.Models; using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; -using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Models; @@ -13,7 +12,7 @@ namespace Elsa.Workflows.Management.Services; public class WorkflowReferenceUpdater( IWorkflowDefinitionPublisher publisher, IWorkflowDefinitionService workflowDefinitionService, - IWorkflowDefinitionStore workflowDefinitionStore, + IWorkflowReferenceQuery workflowReferenceQuery, IApiSerializer serializer) : IWorkflowReferenceUpdater { /// @@ -21,10 +20,17 @@ public class WorkflowReferenceUpdater( { // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. if (referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) - return new UpdateWorkflowReferencesResult([]); + return new([]); - // Find all workflow graphs that contain the updated workflow definition. - var consumingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(referencedDefinition.DefinitionId, cancellationToken); + // Find all workflow definitions that reference the updated workflow definition. + var referencedDefinitionIds = (await workflowReferenceQuery.ExecuteAsync(referencedDefinition.DefinitionId, cancellationToken)).ToList(); + var filter = new WorkflowDefinitionFilter + { + DefinitionIds = referencedDefinitionIds, + VersionOptions = VersionOptions.Latest, + IsReadonly = false + }; + var consumingWorkflowGraphs = await workflowDefinitionService.FindWorkflowGraphsAsync(filter, cancellationToken); // Update consuming workflows. var updatedWorkflows = new List(); @@ -36,7 +42,7 @@ public class WorkflowReferenceUpdater( updatedWorkflows.Add(newDefinition); } - return new UpdateWorkflowReferencesResult(updatedWorkflows); + return new(updatedWorkflows); } private async Task UpdateConsumingWorkflowAsync(WorkflowGraph workflowGraph, WorkflowDefinition definition, CancellationToken cancellationToken) @@ -85,16 +91,6 @@ public class WorkflowReferenceUpdater( .Where(x => x.WorkflowDefinitionId == updatedDefinition.DefinitionId && x.WorkflowDefinitionVersionId != updatedDefinition.Id); } - private async Task> FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(string updatedWorkflowDefinitionId, CancellationToken cancellationToken) - { - var allWorkflowGraphs = (await GetAllWorkflowGraphsAsync(cancellationToken)).ToList(); - var filteredWorkflowGraphs = allWorkflowGraphs - .Where(x => FindWorkflowActivityDefinitionActivityNodes(x.Root).Any(y => y.WorkflowDefinitionId == updatedWorkflowDefinitionId)) - .ToList(); - - return filteredWorkflowGraphs; - } - private IEnumerable FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent) { foreach (var child in parent.Children) @@ -108,32 +104,4 @@ public class WorkflowReferenceUpdater( } } } - - private async Task> GetAllWorkflowGraphsAsync(CancellationToken cancellationToken) - { - var workflowDefinitionFilter = new WorkflowDefinitionFilter - { - VersionOptions = VersionOptions.LatestOrPublished, - IsReadonly = false - }; - var workflowDefinitionSummaries = await workflowDefinitionStore.FindSummariesAsync(workflowDefinitionFilter, cancellationToken); - - // If there are workflow definition summaries with the same definition ID, only take the latest version. - workflowDefinitionSummaries = workflowDefinitionSummaries - .GroupBy(x => x.DefinitionId) - .Select(x => x.OrderByDescending(y => y.Version).First()) - .ToList(); - - var workflowGraphs = new List(); - - foreach (var workflowDefinitionSummary in workflowDefinitionSummaries) - { - var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.Id, cancellationToken); - - if (workflowGraph != null) - workflowGraphs.Add(workflowGraph); - } - - return workflowGraphs; - } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs index 5bd0e7939..96997e284 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -118,7 +118,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP // Serialize materializer context. var materializerContext = materializedWorkflow.MaterializerContext; - var materializerContextJson = materializerContext != null ? _payloadSerializer.Serialize(materializerContext) : default; + var materializerContextJson = materializerContext != null ? _payloadSerializer.Serialize(materializerContext) : null; // Serialize the workflow root. var workflowJson = _activitySerializer.Serialize(workflow.Root);