Merge remote-tracking branch 'origin/blueberry'

This commit is contained in:
Sipke Schoorstra 2024-11-08 21:50:55 +01:00
commit e79b8c87d7
17 changed files with 224 additions and 92 deletions

View file

@ -146,7 +146,6 @@
<PackageVersion Include="Pomelo.EntityFrameworkCore.MySql" Version="7.0.0" />
<PackageVersion Include="Refit" Version="7.1.2" />
<PackageVersion Include="Refit.HttpClientFactory" Version="7.0.0" />
<PackageVersion Include="Scrutor" Version="4.2.2" />
</ItemGroup>
<ItemGroup Condition="'$(TargetFramework)' == 'net8.0'">
<PackageVersion Include="AspNetCore.Authentication.ApiKey" Version="8.0.1" />
@ -182,8 +181,8 @@
<PackageVersion Include="Npgsql.EntityFrameworkCore.PostgreSQL" Version="8.0.10" />
<PackageVersion Include="Polly" Version="8.4.2" />
<PackageVersion Include="Pomelo.EntityFrameworkCore.MySql" Version="8.0.2" />
<PackageVersion Include="Refit" Version="7.2.1" />
<PackageVersion Include="Refit.HttpClientFactory" Version="7.2.1" />
<PackageVersion Include="Scrutor" Version="5.0.1" />
<PackageVersion Include="Microsoft.Extensions.DependencyModel" Version="8.0.2" />
<PackageVersion Include="Refit" Version="8.0.0" />
<PackageVersion Include="Refit.HttpClientFactory" Version="8.0.0" />
</ItemGroup>
</Project>

View file

@ -66,10 +66,8 @@ internal class BulkPublish(IWorkflowDefinitionStore store, IWorkflowDefinitionPu
var result = await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken);
published.Add(definitionId);
if (result.ConsumingWorkflows?.Any() == true)
{
updatedConsumers.AddRange(result.ConsumingWorkflows.Select(x => x.DefinitionId));
}
if (result.AffectedWorkflows.WorkflowDefinitions.Count > 0)
updatedConsumers.AddRange(result.AffectedWorkflows.WorkflowDefinitions.Select(x => x.DefinitionId));
}
return new Response(published, alreadyPublished, notFound, skipped, updatedConsumers);

View file

@ -87,7 +87,7 @@ internal class Post(
PublishWorkflowDefinitionResult? result = null;
if (request.Publish.GetValueOrDefault(false))
if (request.Publish == true)
{
result = await workflowDefinitionPublisher.PublishAsync(draft, cancellationToken);
@ -106,7 +106,8 @@ internal class Post(
}
var mappedDefinition = await linker.MapAsync(draft, cancellationToken);
var response = new Response(mappedDefinition, false, result?.ConsumingWorkflows?.Count() ?? 0);
var affectedWorkflows = result?.AffectedWorkflows?.WorkflowDefinitions ?? [];
var response = new Response(mappedDefinition, false, affectedWorkflows.Count);
await HttpContext.Response.WriteAsJsonAsync(response, serializerOptions, cancellationToken);
}
}

View file

@ -46,7 +46,7 @@ internal class Publish(IWorkflowDefinitionStore store, IWorkflowDefinitionPublis
var isPublished = definition.IsPublished;
var result = !isPublished ? await workflowDefinitionPublisher.PublishAsync(definition, cancellationToken) : null;
var mappedDefinition = await linker.MapAsync(definition, cancellationToken);
var response = new Response(mappedDefinition, isPublished, result?.ConsumingWorkflows?.Count() ?? 0);
var response = new Response(mappedDefinition, isPublished, result?.AffectedWorkflows.WorkflowDefinitions.Count ?? 0);
await SendOkAsync(response, cancellationToken);
}
}

View file

@ -10,7 +10,7 @@ using Microsoft.AspNetCore.Authorization;
namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.UpdateReferences;
[PublicAPI]
internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinitionPublisher workflowDefinitionPublisher, IAuthorizationService authorizationService)
internal class UpdateReferences(IWorkflowReferenceUpdater workflowReferenceUpdater, IWorkflowDefinitionStore store, IAuthorizationService authorizationService)
: ElsaEndpoint<Request, Response>
{
public override void Configure()
@ -43,7 +43,8 @@ internal class UpdateReferences(IWorkflowDefinitionStore store, IWorkflowDefinit
return;
}
var affectedWorkflows = await workflowDefinitionPublisher.UpdateReferencesInConsumingWorkflows(definition, cancellationToken);
var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken);
var affectedWorkflows = result.UpdatedWorkflows;
var response = new Response(affectedWorkflows.Select(w => w.Name ?? w.DefinitionId));
await SendOkAsync(response, cancellationToken);
}

View file

@ -51,7 +51,7 @@ public class WorkflowDefinitionActivity : Composite, IInitializable
var serviceProvider = context.ServiceProvider;
var cancellationToken = context.CancellationToken;
// Find the workflow definition and not the graph; the graph must be computed at runtime, since one NodeIds will vary across graphs.
// Find the workflow definition and not the graph; the graph must be computed at runtime, since NodeIds will vary across graphs.
var workflowDefinition = await GetWorkflowDefinitionAsync(serviceProvider, cancellationToken);
if (workflowDefinition == null)

View file

@ -67,12 +67,4 @@ public interface IWorkflowDefinitionPublisher
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The saved workflow definition.</returns>
Task<WorkflowDefinition> SaveDraftAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default);
/// <summary>
/// Updates all referencing workflow definitions to use the version of the specified workflow definition.
/// </summary>
/// <param name="dependency">The workflow definition to update references for.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The updated workflow definitions.</returns>
Task<IEnumerable<WorkflowDefinition>> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default);
}

View file

@ -0,0 +1,18 @@
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Models;
namespace Elsa.Workflows.Management.Contracts;
/// <summary>
/// Updates references to the specified workflow of all workflows that reference it.
/// </summary>
public interface IWorkflowReferenceUpdater
{
/// <summary>
/// Updates references to the specified workflow of all workflows that reference it.
/// </summary>
/// <param name="referencedDefinition">The workflow definition that is being referenced. All workflows that reference this definition will be updated to use this newest version.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The result of the operation.</returns>
Task<UpdateWorkflowReferencesResult> UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default);
}

View file

@ -215,6 +215,7 @@ public class WorkflowManagementFeature : FeatureBase
.AddScoped<IWorkflowDefinitionImporter, WorkflowDefinitionImporter>()
.AddScoped<IWorkflowDefinitionManager, WorkflowDefinitionManager>()
.AddScoped<IWorkflowInstanceManager, WorkflowInstanceManager>()
.AddScoped<IWorkflowReferenceUpdater, WorkflowReferenceUpdater>()
.AddScoped<IActivityRegistryPopulator, ActivityRegistryPopulator>()
.AddSingleton<IExpressionDescriptorRegistry, ExpressionDescriptorRegistry>()
.AddSingleton<IExpressionDescriptorProvider, DefaultExpressionDescriptorProvider>()
@ -235,6 +236,7 @@ public class WorkflowManagementFeature : FeatureBase
Services
.AddNotificationHandler<DeleteWorkflowInstances>()
.AddNotificationHandler<RefreshActivityRegistry>()
.AddNotificationHandler<UpdateConsumingWorkflows>()
;
Services.Configure<ManagementOptions>(options =>

View file

@ -0,0 +1,27 @@
using Elsa.Common.Models;
using Elsa.Extensions;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
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.Notifications;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Handlers;
/// <summary>
/// Updates consuming workflows when a workflow definition is published.
/// </summary>
public class UpdateConsumingWorkflows(IWorkflowReferenceUpdater workflowReferenceUpdater) : INotificationHandler<WorkflowDefinitionPublished>
{
/// <inheritdoc />
public async Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken)
{
var definition = notification.WorkflowDefinition;
var result = await workflowReferenceUpdater.UpdateWorkflowReferencesAsync(definition, cancellationToken);
notification.AffectedWorkflows.WorkflowDefinitions.AddRange(result.UpdatedWorkflows);
}
}

View file

@ -0,0 +1,5 @@
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Management.Models;
public record AffectedWorkflows(ICollection<WorkflowDefinition> WorkflowDefinitions);

View file

@ -1,8 +1,6 @@
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Management.Models;
/// <summary>
/// Represents the result of publishing a workflow definition.
/// </summary>
public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection<WorkflowValidationError> ValidationErrors, IEnumerable<WorkflowDefinition>? ConsumingWorkflows);
public record PublishWorkflowDefinitionResult(bool Succeeded, ICollection<WorkflowValidationError> ValidationErrors, AffectedWorkflows AffectedWorkflows);

View file

@ -0,0 +1,18 @@
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Models;
/// <summary>
/// Represents the result of updating workflow references.
/// </summary>
public class UpdateWorkflowReferencesResult(IEnumerable<WorkflowDefinition> updatedWorkflows)
{
/// <summary>
/// Gets a collection of workflow graphs that have been updated.
/// </summary>
/// <value>
/// A read-only collection of <see cref="WorkflowGraph"/> instances representing the updated workflows.
/// </value>
public IReadOnlyCollection<WorkflowDefinition> UpdatedWorkflows { get; } = updatedWorkflows.ToList();
}

View file

@ -1,5 +1,6 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Models;
using JetBrains.Annotations;
namespace Elsa.Workflows.Management.Notifications;
@ -9,4 +10,4 @@ namespace Elsa.Workflows.Management.Notifications;
/// </summary>
/// <param name="WorkflowDefinition">The workflow definition.</param>
[PublicAPI]
public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition) : INotification;
public record WorkflowDefinitionPublished(WorkflowDefinition WorkflowDefinition, AffectedWorkflows AffectedWorkflows) : INotification;

View file

@ -1,7 +1,6 @@
using Elsa.Common;
using Elsa.Common.Entities;
using Elsa.Common.Models;
using Elsa.Extensions;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
@ -21,7 +20,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
private readonly IWorkflowDefinitionStore _workflowDefinitionStore;
private readonly INotificationSender _notificationSender;
private readonly IIdentityGenerator _identityGenerator;
private readonly IActivityVisitor _activityVisitor;
private readonly IActivitySerializer _activitySerializer;
private readonly IRequestSender _requestSender;
private readonly ISystemClock _systemClock;
@ -34,7 +32,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
IWorkflowDefinitionStore workflowDefinitionStore,
INotificationSender notificationSender,
IIdentityGenerator identityGenerator,
IActivityVisitor activityVisitor,
IActivitySerializer activitySerializer,
IRequestSender requestSender,
ISystemClock systemClock)
@ -43,7 +40,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
_workflowDefinitionStore = workflowDefinitionStore;
_notificationSender = notificationSender;
_identityGenerator = identityGenerator;
_activityVisitor = activityVisitor;
_activitySerializer = activitySerializer;
_requestSender = requestSender;
_systemClock = systemClock;
@ -119,16 +115,9 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
definition = Initialize(definition);
await _workflowDefinitionStore.SaveAsync(definition, cancellationToken);
await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition), cancellationToken);
var consumingWorkflows = new List<WorkflowDefinition>();
if (definition.Options is { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true })
{
consumingWorkflows.AddRange(await UpdateReferencesInConsumingWorkflows(definition, cancellationToken));
}
return new PublishWorkflowDefinitionResult(true, validationErrors, consumingWorkflows);
var affectedWorkflows = new AffectedWorkflows(new List<WorkflowDefinition>());
await _notificationSender.SendAsync(new WorkflowDefinitionPublished(definition, affectedWorkflows), cancellationToken);
return new PublishWorkflowDefinitionResult(true, validationErrors, affectedWorkflows);
}
/// <inheritdoc />
@ -210,57 +199,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
return draft;
}
/// <inheritdoc />
public async Task<IEnumerable<WorkflowDefinition>> UpdateReferencesInConsumingWorkflows(WorkflowDefinition dependency, CancellationToken cancellationToken = default)
{
var updatedWorkflowDefinitions = new List<WorkflowDefinition>();
var workflowDefinitions = (await _workflowDefinitionStore.FindManyAsync(new WorkflowDefinitionFilter
{
VersionOptions = VersionOptions.LatestOrPublished
}, cancellationToken)).ToList();
// Remove the dependency from the list of workflow definitions to consider.
workflowDefinitions = workflowDefinitions.Where(x => x.DefinitionId != dependency.DefinitionId).ToList();
foreach (var definition in workflowDefinitions)
{
var root = _activitySerializer.Deserialize(definition.StringData!);
var graph = await _activityVisitor.VisitAsync(root, cancellationToken);
var flattenedList = graph.Flatten().ToList();
var definitionId = dependency.DefinitionId;
var version = dependency.Version;
var nodes = flattenedList
.Where(x => x.Activity is WorkflowDefinitionActivity workflowDefinitionActivity && workflowDefinitionActivity.WorkflowDefinitionId == definitionId)
.ToList();
foreach (var node in nodes.Where(activity => activity.Activity.Version < version))
{
var activity = (WorkflowDefinitionActivity)node.Activity;
activity.Version = version;
activity.WorkflowDefinitionVersionId = dependency.Id;
if (!updatedWorkflowDefinitions.Contains(definition))
updatedWorkflowDefinitions.Add(definition);
}
if (updatedWorkflowDefinitions.Contains(definition))
{
var serializedData = _activitySerializer.Serialize(root);
definition.StringData = serializedData;
}
}
if (updatedWorkflowDefinitions.Any())
{
await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdating(updatedWorkflowDefinitions), cancellationToken);
await _workflowDefinitionStore.SaveManyAsync(updatedWorkflowDefinitions, cancellationToken);
await _notificationSender.SendAsync(new WorkflowDefinitionVersionsUpdated(updatedWorkflowDefinitions), cancellationToken);
}
return updatedWorkflowDefinitions;
}
private WorkflowDefinition Initialize(WorkflowDefinition definition)
{
if (definition.Id == null!)

View file

@ -0,0 +1,131 @@
using Elsa.Common.Models;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Contracts;
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;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Services;
/// <inheritdoc />
public class WorkflowReferenceUpdater(
IWorkflowDefinitionPublisher publisher,
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowDefinitionStore workflowDefinitionStore,
IApiSerializer serializer) : IWorkflowReferenceUpdater
{
/// <inheritdoc />
public async Task<UpdateWorkflowReferencesResult> UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default)
{
// 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([]);
// Find all workflow graphs that contain the updated workflow definition.
var consumingWorkflowGraphs = await FindWorkflowsContainingUpdatedWorkflowDefinitionAsync(referencedDefinition.DefinitionId, cancellationToken);
// Update consuming workflows.
var updatedWorkflows = new List<WorkflowDefinition>();
foreach (var workflowGraph in consumingWorkflowGraphs)
{
var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, referencedDefinition, cancellationToken);
if (newDefinition != null)
updatedWorkflows.Add(newDefinition);
}
return new UpdateWorkflowReferencesResult(updatedWorkflows);
}
private async Task<WorkflowDefinition?> UpdateConsumingWorkflowAsync(WorkflowGraph workflowGraph, WorkflowDefinition definition, CancellationToken cancellationToken)
{
// Create a new version of the published workflow definition or get the existing draft.
var consumerDefinitionId = workflowGraph.Workflow.Identity.DefinitionId;
var originalVersionIsPublished = workflowGraph.Workflow.Publication.IsPublished;
var newVersion = await publisher.GetDraftAsync(consumerDefinitionId, VersionOptions.LatestOrPublished, cancellationToken);
// This is null in case the definition no longer exists in the store.
if (newVersion == null)
return null;
// Materialize the draft to find all workflow definition activities that use the updated workflow definition.
var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken);
var outdatedWorkflowDefinitionActivities = FindOutdatedWorkflowDefinitionActivities(newWorkflowGraph, definition).ToList();
// Skip if the new version of the published workflow definition is not used in the workflow or if the activity is already up to date.
if (outdatedWorkflowDefinitionActivities.Count == 0)
return null;
// Update the consuming workflow graph to use the new version of the published workflow definition.
foreach (var workflowDefinitionActivity in outdatedWorkflowDefinitionActivities)
{
workflowDefinitionActivity.WorkflowDefinitionVersionId = definition.Id;
workflowDefinitionActivity.LatestAvailablePublishedVersionId = definition.Id;
workflowDefinitionActivity.Version = definition.Version;
}
// Update the new version of the published workflow definition.
if (newWorkflowGraph.Root.Activity is Workflow newWorkflow)
newVersion.StringData = serializer.Serialize(newWorkflow.Root);
// If the draft is new, publish it.
if (originalVersionIsPublished)
await publisher.PublishAsync(newVersion, cancellationToken);
else
await publisher.SaveDraftAsync(newVersion, cancellationToken);
return newVersion;
}
private IEnumerable<WorkflowDefinitionActivity> FindOutdatedWorkflowDefinitionActivities(WorkflowGraph workflowGraph, WorkflowDefinition updatedDefinition)
{
return FindWorkflowActivityDefinitionActivityNodes(workflowGraph.Root)
.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)
{
if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity)
yield return workflowDefinitionActivity;
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);
var workflowGraphs = new List<WorkflowGraph>();
foreach (var workflowDefinitionSummary in workflowDefinitionSummaries)
{
var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionSummary.DefinitionId, VersionOptions.LatestOrPublished, cancellationToken);
if (workflowGraph != null)
workflowGraphs.Add(workflowGraph);
}
return workflowGraphs;
}
}

View file

@ -35,7 +35,7 @@ public class ReloadWorkflowTests : AppComponentTest
var client = WorkflowServer.CreateHttpWorkflowClient();
await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None);
var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test"));
await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None);
await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync();
var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test"));
Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode);
Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode);
@ -92,6 +92,9 @@ public class ReloadWorkflowTests : AppComponentTest
// Assert that the activity registry contains a new activity descriptor representing the new workflow version.
var activityV2 = _activityRegistry.Find(activityTypeName)!;
Assert.Equal(2, activityV2.Version);
// Cleanup: Delete the workflow definition and its versions.
await _workflowDefinitionManager.DeleteByDefinitionIdAsync(definitionId, CancellationToken.None);
}
private async Task<MaterializedWorkflow> BuildWorkflowAsync(string definitionId, string definitionVersionId, int version)