Remove UpdateReferencesInConsumingWorkflows method
Eliminated the UpdateReferencesInConsumingWorkflows method from IWorkflowDefinitionPublisher and its implementation from WorkflowDefinitionPublisher. The responsibility of updating workflow references is now transitioned to IWorkflowReferenceUpdater used in the `UpdateReferences` endpoint.
This commit is contained in:
parent
a7b0d5630d
commit
8fe0147abd
|
|
@ -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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -68,12 +68,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);
|
||||
}
|
||||
|
|
@ -1,11 +1,9 @@
|
|||
using Elsa.Common.Contracts;
|
||||
using Elsa.Common.Entities;
|
||||
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;
|
||||
|
|
@ -23,7 +21,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;
|
||||
|
|
@ -36,7 +33,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
|
|||
IWorkflowDefinitionStore workflowDefinitionStore,
|
||||
INotificationSender notificationSender,
|
||||
IIdentityGenerator identityGenerator,
|
||||
IActivityVisitor activityVisitor,
|
||||
IActivitySerializer activitySerializer,
|
||||
IRequestSender requestSender,
|
||||
ISystemClock systemClock)
|
||||
|
|
@ -45,7 +41,6 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher
|
|||
_workflowDefinitionStore = workflowDefinitionStore;
|
||||
_notificationSender = notificationSender;
|
||||
_identityGenerator = identityGenerator;
|
||||
_activityVisitor = activityVisitor;
|
||||
_activitySerializer = activitySerializer;
|
||||
_requestSender = requestSender;
|
||||
_systemClock = systemClock;
|
||||
|
|
@ -205,57 +200,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!)
|
||||
|
|
|
|||
Loading…
Reference in a new issue