Refactor database extensions and support migrations for V3.6 (#6788)

* Refactor database extensions and support migrations for V3.6

Remove `DatabaseFacadeExtensions` and introduce `IWorkflowReferenceQuery` with its default implementation. Implement database schema updates for PostgreSQL, MySQL, and Oracle to enhance compatibility with the V3.6 data structure.

* Remove commented-out code and standardize null default assignment in `IWorkflowDefinitionStore` interface
This commit is contained in:
Sipke Schoorstra 2025-07-14 09:48:22 +02:00 committed by GitHub
parent 0edc744172
commit cfb48ffcbf
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
9 changed files with 151 additions and 56 deletions

View file

@ -55,4 +55,9 @@ public interface IWorkflowDefinitionService
/// Looks for a <see cref="WorkflowGraph"/> by the specified <see cref="WorkflowDefinitionFilter"/>.
/// </summary>
Task<WorkflowGraph?> FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
/// <summary>
/// Looks for all <see cref="WorkflowGraph"/>s that match the specified <see cref="WorkflowDefinitionFilter"/>.
/// </summary>
Task<IEnumerable<WorkflowGraph>> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default);
}

View file

@ -159,5 +159,5 @@ public interface IWorkflowDefinitionStore
/// <param name="name">The name.</param>
/// <param name="definitionId">The definition ID to exclude from the check.</param>
/// <param name="cancellationToken">The cancellation token.</param>
Task<bool> GetIsNameUnique(string name, string? definitionId = default, CancellationToken cancellationToken = default);
Task<bool> GetIsNameUnique(string name, string? definitionId = null, CancellationToken cancellationToken = default);
}

View file

@ -0,0 +1,15 @@
namespace Elsa.Workflows.Management;
/// <summary>
/// Finds all latest versions of workflow definitions that reference a specific workflow definition.
/// </summary>
public interface IWorkflowReferenceQuery
{
/// <summary>
/// Queries all latest versions of workflow definitions that reference the specified workflow definition.
/// </summary>
/// <param name="workflowDefinitionId">The ID of the workflow definition to query references for.</param>
/// <param name="cancellationToken">The cancellation token to cancel the operation.</param>
/// <returns>A collection of workflow definition IDs that reference the specified workflow definition.</returns>
Task<IEnumerable<string>> ExecuteAsync(string workflowDefinitionId, CancellationToken cancellationToken = default);
}

View file

@ -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<IServiceProvider, IWorkflowDefinitionPublisher> _workflowDefinitionPublisher = sp => ActivatorUtilities.CreateInstance<WorkflowDefinitionPublisher>(sp);
private Func<IServiceProvider, IWorkflowReferenceQuery> _workflowReferenceQuery = sp => ActivatorUtilities.CreateInstance<DefaultWorkflowReferenceQuery>(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<IServiceProvider, IWorkflowDefinitionPublisher> workflowDefinitionPublisher)
public WorkflowManagementFeature UseWorkflowDefinitionPublisher(Func<IServiceProvider, IWorkflowDefinitionPublisher> workflowDefinitionPublisher)
{
_workflowDefinitionPublisher = workflowDefinitionPublisher;
return this;
}
public WorkflowManagementFeature UseWorkflowReferenceFinder<T>() where T : class, IWorkflowReferenceQuery
{
Services.TryAddScoped<T>();
return UseWorkflowReferenceFinder(sp => sp.GetRequiredService<T>());
}
public WorkflowManagementFeature UseWorkflowReferenceFinder(Func<IServiceProvider, IWorkflowReferenceQuery> workflowReferenceFinder)
{
_workflowReferenceQuery = workflowReferenceFinder;
return this;
}
/// <inheritdoc />
[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<IWorkflowDefinitionService, WorkflowDefinitionService>()
.AddScoped<IWorkflowSerializer, WorkflowSerializer>()
.AddScoped<IWorkflowValidator, WorkflowValidator>()
.AddScoped(_workflowReferenceQuery)
.AddScoped(_workflowDefinitionPublisher)
.AddScoped<IWorkflowDefinitionImporter, WorkflowDefinitionImporter>()
.AddScoped<IWorkflowDefinitionManager, WorkflowDefinitionManager>()

View file

@ -11,7 +11,7 @@ namespace Elsa.Workflows.Management.Services;
/// Decorates an <see cref="IWorkflowDefinitionService"/> with caching capabilities.
/// </summary>
[UsedImplicitly]
public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager) : IWorkflowDefinitionService
public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowDefinitionService
{
/// <inheritdoc />
public async Task<WorkflowGraph> MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
@ -41,10 +41,7 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat
/// <inheritdoc />
public Task<WorkflowDefinition?> 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
/// <inheritdoc />
public Task<WorkflowGraph?> 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<IEnumerable<WorkflowGraph>> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken);
var workflowGraphs = new List<WorkflowGraph>();
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<T?> GetFromCacheAsync<T>(string cacheKey, Func<Task<T?>> getObjectFunc, Func<T, string> getChangeTokenKeyFunc)
{
var cache = cacheManager.Cache;

View file

@ -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<IEnumerable<string>> ExecuteAsync(string workflowDefinitionId, CancellationToken cancellationToken = default)
{
var workflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(workflowDefinitionId, cancellationToken);
return workflowGraphs.Select(x => x.Workflow.Identity.DefinitionId).Distinct();
}
private async Task<IEnumerable<WorkflowGraph>> 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<WorkflowDefinitionActivity> 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<IEnumerable<WorkflowGraph>> 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<WorkflowGraph>();
foreach (var workflowDefinitionSummary in workflowDefinitionSummaries)
{
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.Id, cancellationToken);
if (workflowGraph != null)
workflowGraphs.Add(workflowGraph);
}
return workflowGraphs;
}
}

View file

@ -95,4 +95,18 @@ public class WorkflowDefinitionService(
return await MaterializeWorkflowAsync(definition, cancellationToken);
}
/// <inheritdoc />
public async Task<IEnumerable<WorkflowGraph>> FindWorkflowGraphsAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
{
var workflowDefinitions = await workflowDefinitionStore.FindManyAsync(filter, cancellationToken);
var workflowGraphs = new List<WorkflowGraph>();
foreach (var workflowDefinition in workflowDefinitions)
{
var workflowGraph = await MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
workflowGraphs.Add(workflowGraph);
}
return workflowGraphs;
}
}

View file

@ -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
{
/// <inheritdoc />
@ -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<WorkflowDefinition>();
@ -36,7 +42,7 @@ public class WorkflowReferenceUpdater(
updatedWorkflows.Add(newDefinition);
}
return new UpdateWorkflowReferencesResult(updatedWorkflows);
return new(updatedWorkflows);
}
private async Task<WorkflowDefinition?> 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<IEnumerable<WorkflowGraph>> 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<WorkflowDefinitionActivity> FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent)
{
foreach (var child in parent.Children)
@ -108,32 +104,4 @@ public class WorkflowReferenceUpdater(
}
}
}
private async Task<IEnumerable<WorkflowGraph>> 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<WorkflowGraph>();
foreach (var workflowDefinitionSummary in workflowDefinitionSummaries)
{
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.Id, cancellationToken);
if (workflowGraph != null)
workflowGraphs.Add(workflowGraph);
}
return workflowGraphs;
}
}

View file

@ -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);