Refactors workflow reference updates (#6792)

* 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

* Add XML documentation for `DefaultWorkflowReferenceQuery` detailing its purpose and dependencies

* Refactor `WorkflowReferenceUpdater` to support recursive dependency resolution, prevent concurrent updates, and improve reference consistency.

* Simplify `WorkflowReferenceUpdater` by removing topological sorting and redundant dependencies handling.

* Refactor `WorkflowReferenceUpdater` to streamline reference updates, remove redundant logic, and enhance dependency resolution efficiency.

* Refactor `WorkflowReferenceUpdater` to use `HashSet` for updated workflows, reducing potential duplication and improving performance.

* Refactor `WorkflowReferenceUpdater` to introduce topological sorting for correct processing order, improve dependency resolution, and enhance clarity with updated records and comments.

* Introduce `WorkflowDefinitionActivityDescriptorFactory` to simplify `WorkflowDefinitionActivity` descriptor creation and refactor existing components for modularity, clarity, and efficiency.

* Update `WorkflowReferenceUpdater` to use `VersionOptions.Latest` instead of `VersionOptions.LatestOrPublished` for workflow reference resolution.

* Refactor `WorkflowReferenceUpdater` to improve workflow dependency resolution by handling publication states, caching drafts more efficiently, and introducing distinct processing for latest and published versions.

* Refactor workflow publication logic and update SQLite configuration.

Removed unused draft publication logic to simplify workflow reference updates. Updated SQLite persistence configuration in `Elsa.Server.Agents.Web` to use explicit connection strings for improved clarity and maintainability.

* Remove commented-out legacy code in `WorkflowReferenceUpdater` to improve clarity and maintainability.

* Fix formatting by adding a missing newline at EOF in `Directory.Build.props`.

* Prevent infinite recursion in `GetReferencingWorkflowDefinitionIdsAsync` by introducing visited ID tracking. Fix formatting inconsistencies in `WorkflowReferenceUpdater`.

* Update `WorkflowReferenceUpdater` to use `NewGraph` instead of materializing workflows for referencing workflow graphs
This commit is contained in:
Sipke Schoorstra 2025-07-15 13:59:05 +02:00 committed by GitHub
parent cfb48ffcbf
commit 3f5cac76c5
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
7 changed files with 352 additions and 216 deletions

View file

@ -34,9 +34,9 @@ public class RawStringContent : HttpContent
/// <inheritdoc />
protected override async Task SerializeToStreamAsync(Stream stream, TransportContext? context, CancellationToken cancellationToken)
{
using var writer = new StreamWriter(stream, _encoding, leaveOpen: true);
await using var writer = new StreamWriter(stream, _encoding, leaveOpen: true);
await writer.WriteAsync(_content.AsMemory(), cancellationToken);
await writer.FlushAsync();
await writer.FlushAsync(cancellationToken);
}
/// <inheritdoc />

View file

@ -0,0 +1,120 @@
using Elsa.Extensions;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
using Humanizer;
namespace Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
public class WorkflowDefinitionActivityDescriptorFactory(IActivityFactory activityFactory)
{
public ActivityDescriptor CreateDescriptor(WorkflowDefinition definition, WorkflowDefinition? latestPublishedDefinition = null)
{
var typeName = definition.Name!.Pascalize();
var ports = definition.Outcomes.Select(outcome => new Port
{
Name = outcome,
DisplayName = outcome,
IsBrowsable = true,
Type = PortType.Flow
}).ToList();
var rootPort = new Port
{
Name = nameof(WorkflowDefinitionActivity.Root),
DisplayName = "Root",
IsBrowsable = false,
Type = PortType.Embedded
};
ports.Insert(0, rootPort);
return new()
{
TypeName = typeName,
Name = typeName,
Version = definition.Version,
DisplayName = definition.Name,
Description = definition.Description,
Category = definition.Options.ActivityCategory ?? "Workflows",
Kind = ActivityKind.Action,
IsBrowsable = definition.IsPublished,
Inputs = DescribeInputs(definition).ToList(),
Outputs = DescribeOutputs(definition).ToList(),
Ports = ports,
CustomProperties =
{
["RootType"] = nameof(WorkflowDefinitionActivity),
["WorkflowDefinitionId"] = definition.DefinitionId,
["WorkflowDefinitionVersionId"] = definition.Id
},
ConstructionProperties = new Dictionary<string, object>
{
[nameof(WorkflowDefinitionActivity.WorkflowDefinitionId)] = definition.DefinitionId,
[nameof(WorkflowDefinitionActivity.WorkflowDefinitionVersionId)] = definition.Id,
[nameof(WorkflowDefinitionActivity.Version)] = definition.Version,
},
Constructor = context =>
{
var activity = (WorkflowDefinitionActivity)activityFactory.Create(typeof(WorkflowDefinitionActivity), context);
activity.Type = typeName;
activity.WorkflowDefinitionId = definition.DefinitionId;
activity.WorkflowDefinitionVersionId = definition.Id;
activity.Version = definition.Version;
activity.LatestAvailablePublishedVersion = latestPublishedDefinition?.Version ?? definition.Version;
activity.LatestAvailablePublishedVersionId = latestPublishedDefinition?.Id ?? definition.Id;
return activity;
}
};
}
private static IEnumerable<InputDescriptor> DescribeInputs(WorkflowDefinition definition)
{
var inputs = definition.Inputs.Select(inputDefinition =>
{
var nakedType = inputDefinition.Type;
var inputName = inputDefinition.Name;
var safeInputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), inputName);
return new InputDescriptor
{
Type = nakedType,
IsWrapped = true,
ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeInputName),
ValueSetter = (activity, value) => activity.SyntheticProperties[safeInputName] = value!,
Name = safeInputName,
DisplayName = inputDefinition.DisplayName,
Description = inputDefinition.Description,
Category = inputDefinition.Category,
UIHint = inputDefinition.UIHint,
StorageDriverType = inputDefinition.StorageDriverType,
IsSynthetic = true
};
});
foreach (var input in inputs)
yield return input;
}
private static IEnumerable<OutputDescriptor> DescribeOutputs(WorkflowDefinition definition)
{
return definition.Outputs.Select(outputDefinition =>
{
var nakedType = outputDefinition.Type;
var outputName = outputDefinition.Name;
var safeOutputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), outputName);
return new OutputDescriptor
{
Type = nakedType,
ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeOutputName),
ValueSetter = (activity, value) => activity.SyntheticProperties[safeOutputName] = value!,
Name = safeOutputName,
DisplayName = outputDefinition.DisplayName,
Description = outputDefinition.Description,
IsSynthetic = true
};
});
}
}

View file

@ -1,18 +1,14 @@
using Elsa.Common.Models;
using Elsa.Extensions;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Models;
using Elsa.Workflows.Serialization.Converters;
using Elsa.Workflows.Serialization.Helpers;
using Humanizer;
namespace Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
/// <summary>
/// Provides activity descriptors based on <see cref="WorkflowDefinition"/>s stored in the database.
/// </summary>
public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store, IActivityFactory activityFactory, ActivityWriter activityWriter) : IActivityProvider
public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store, WorkflowDefinitionActivityDescriptorFactory workflowDefinitionActivityDescriptorFactory) : IActivityProvider
{
/// <inheritdoc />
public async ValueTask<IEnumerable<ActivityDescriptor>> GetDescriptorsAsync(CancellationToken cancellationToken = default)
@ -22,24 +18,9 @@ public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store,
UsableAsActivity = true,
VersionOptions = VersionOptions.All
};
var allDescriptors = new List<ActivityDescriptor>();
var currentPage = 0;
const int pageSize = 100;
while (true)
{
var pageArgs = PageArgs.FromPage(currentPage++, pageSize);
var pageOfDefinitions = await store.FindManyAsync(filter, pageArgs, cancellationToken);
var descriptors = CreateDescriptors(pageOfDefinitions.Items).ToList();
allDescriptors.AddRange(descriptors);
if (allDescriptors.Count >= pageOfDefinitions.TotalCount)
break;
}
return allDescriptors;
var definitions = (await store.FindManyAsync(filter, cancellationToken)).ToList();
return CreateDescriptors(definitions).ToList();
}
private IEnumerable<ActivityDescriptor> CreateDescriptors(ICollection<WorkflowDefinition> definitions)
@ -49,116 +30,9 @@ public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store,
private ActivityDescriptor CreateDescriptor(WorkflowDefinition definition, ICollection<WorkflowDefinition> allDefinitions)
{
var typeName = definition.Name!.Pascalize();
var latestPublishedVersion = allDefinitions
.Where(x => x.DefinitionId == definition.DefinitionId && x.IsPublished)
.MaxBy(x => x.Version);
var ports = definition.Outcomes.Select(outcome => new Port
{
Name = outcome,
DisplayName = outcome,
IsBrowsable = true,
Type = PortType.Flow
}).ToList();
var rootPort = new Port
{
Name = nameof(WorkflowDefinitionActivity.Root),
DisplayName = "Root",
IsBrowsable = false,
Type = PortType.Embedded
};
ports.Insert(0, rootPort);
return new()
{
TypeName = typeName,
Name = typeName,
Version = definition.Version,
DisplayName = definition.Name,
Description = definition.Description,
Category = definition.Options.ActivityCategory ?? "Workflows",
Kind = ActivityKind.Action,
IsBrowsable = definition.IsPublished,
Inputs = DescribeInputs(definition).ToList(),
Outputs = DescribeOutputs(definition).ToList(),
Ports = ports,
CustomProperties =
{
["RootType"] = nameof(WorkflowDefinitionActivity),
["WorkflowDefinitionId"] = definition.DefinitionId,
["WorkflowDefinitionVersionId"] = definition.Id
},
ConstructionProperties = new Dictionary<string, object>
{
[nameof(WorkflowDefinitionActivity.WorkflowDefinitionId)] = definition.DefinitionId,
[nameof(WorkflowDefinitionActivity.WorkflowDefinitionVersionId)] = definition.Id,
[nameof(WorkflowDefinitionActivity.Version)] = definition.Version,
},
Constructor = context =>
{
var activity = (WorkflowDefinitionActivity)activityFactory.Create(typeof(WorkflowDefinitionActivity), context);
activity.Type = typeName;
activity.WorkflowDefinitionId = definition.DefinitionId;
activity.WorkflowDefinitionVersionId = definition.Id;
activity.Version = definition.Version;
activity.LatestAvailablePublishedVersion = latestPublishedVersion?.Version ?? 0;
activity.LatestAvailablePublishedVersionId = latestPublishedVersion?.Id;
return activity;
}
};
}
private static IEnumerable<InputDescriptor> DescribeInputs(WorkflowDefinition definition)
{
var inputs = definition.Inputs.Select(inputDefinition =>
{
var nakedType = inputDefinition.Type;
var inputName = inputDefinition.Name;
var safeInputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), inputName);
return new InputDescriptor
{
Type = nakedType,
IsWrapped = true,
ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeInputName),
ValueSetter = (activity, value) => activity.SyntheticProperties[safeInputName] = value!,
Name = safeInputName,
DisplayName = inputDefinition.DisplayName,
Description = inputDefinition.Description,
Category = inputDefinition.Category,
UIHint = inputDefinition.UIHint,
StorageDriverType = inputDefinition.StorageDriverType,
IsSynthetic = true
};
});
foreach (var input in inputs)
yield return input;
}
private static IEnumerable<OutputDescriptor> DescribeOutputs(WorkflowDefinition definition)
{
return definition.Outputs.Select(outputDefinition =>
{
var nakedType = outputDefinition.Type;
var outputName = outputDefinition.Name;
var safeOutputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), outputName);
return new OutputDescriptor
{
Type = nakedType,
ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeOutputName),
ValueSetter = (activity, value) => activity.SyntheticProperties[safeOutputName] = value!,
Name = safeOutputName,
DisplayName = outputDefinition.DisplayName,
Description = outputDefinition.Description,
IsSynthetic = true
};
});
return workflowDefinitionActivityDescriptorFactory.CreateDescriptor(definition, latestPublishedVersion);
}
}

View file

@ -226,6 +226,7 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module)
.AddMemoryStore<WorkflowInstance, MemoryWorkflowInstanceStore>()
.AddActivityProvider<TypedActivityProvider>()
.AddActivityProvider<WorkflowDefinitionActivityProvider>()
.AddScoped<WorkflowDefinitionActivityDescriptorFactory>()
.AddScoped<WorkflowDefinitionActivityProvider>()
.AddScoped<IWorkflowDefinitionActivityRegistryUpdater, WorkflowDefinitionActivityRegistryUpdater>()
.AddScoped<IWorkflowDefinitionService, WorkflowDefinitionService>()

View file

@ -5,6 +5,15 @@ using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Services;
/// <summary>
/// Represents the default implementation of <see cref="IWorkflowReferenceQuery"/> that queries workflows
/// referencing a specific workflow definition.
/// </summary>
/// <remarks>
/// This class is designed to identify workflows that are dependent on a given workflow definition.
/// It leverages services such as <see cref="IWorkflowDefinitionService"/> and <see cref="IWorkflowDefinitionStore"/>
/// to locate, graph, and analyze workflows that include the specified workflow definition.
/// </remarks>
public class DefaultWorkflowReferenceQuery(IWorkflowDefinitionService workflowDefinitionService, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowReferenceQuery
{
public async Task<IEnumerable<string>> ExecuteAsync(string workflowDefinitionId, CancellationToken cancellationToken = default)

View file

@ -1,107 +1,229 @@
using System.Runtime.CompilerServices;
using Elsa.Common.Models;
using Elsa.Extensions;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
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 />
internal record WorkflowReferences(string ReferencedDefinitionId, ICollection<string> ReferencingDefinitionIds);
internal record UpdatedWorkflowDefinition(WorkflowDefinition Definition, WorkflowGraph NewGraph);
public class WorkflowReferenceUpdater(
IWorkflowDefinitionPublisher publisher,
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowDefinitionStore workflowDefinitionStore,
IWorkflowReferenceQuery workflowReferenceQuery,
IApiSerializer serializer) : IWorkflowReferenceUpdater
WorkflowDefinitionActivityDescriptorFactory workflowDefinitionActivityDescriptorFactory,
IActivityRegistry activityRegistry,
IApiSerializer serializer)
: IWorkflowReferenceUpdater
{
/// <inheritdoc />
public async Task<UpdateWorkflowReferencesResult> UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default)
private bool _isUpdating;
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 })
if (_isUpdating ||
referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true })
return new([]);
// 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);
var allWorkflowReferences = await GetReferencingWorkflowDefinitionIdsAsync(referencedDefinition.DefinitionId, cancellationToken).ToListAsync(cancellationToken);
var filteredWorkflowReferences = allWorkflowReferences
.Where(r => r.ReferencingDefinitionIds.Any())
.DistinctBy(r => r.ReferencedDefinitionId)
.ToList();
// Update consuming workflows.
var updatedWorkflows = new List<WorkflowDefinition>();
foreach (var workflowGraph in consumingWorkflowGraphs)
{
var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, referencedDefinition, cancellationToken);
var referencingIds = filteredWorkflowReferences.SelectMany(r => r.ReferencingDefinitionIds).Distinct().ToList();
var referencedIds = filteredWorkflowReferences.Select(r => r.ReferencedDefinitionId).Distinct().ToList();
if (newDefinition != null)
updatedWorkflows.Add(newDefinition);
}
return new(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 IEnumerable<WorkflowDefinitionActivity> FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent)
{
foreach (var child in parent.Children)
{
if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity)
yield return workflowDefinitionActivity;
else
var referencingWorkflowGraphs = (await workflowDefinitionService.FindWorkflowGraphsAsync(new()
{
foreach (var grandChild in FindWorkflowActivityDefinitionActivityNodes(child))
yield return grandChild;
DefinitionIds = referencingIds,
VersionOptions = VersionOptions.Latest,
IsReadonly = false
}, cancellationToken))
.ToDictionary(g => g.Workflow.Identity.DefinitionId);
var referencedWorkflowDefinitionList = (await workflowDefinitionStore.FindManyAsync(new()
{
DefinitionIds = referencedIds,
VersionOptions = VersionOptions.Published,
IsReadonly = false
}, cancellationToken)).ToList();
var referencedWorkflowDefinitionsPublished = referencedWorkflowDefinitionList
.GroupBy(x => x.DefinitionId)
.Select(group =>
{
var publishedVersion = group.FirstOrDefault(x => x.IsPublished);
return publishedVersion ?? group.First();
})
.ToDictionary(d => d.DefinitionId);
var initialPublicationState = new Dictionary<string, bool>();
foreach (var workflowGraph in referencingWorkflowGraphs)
initialPublicationState[workflowGraph.Key] = workflowGraph.Value.Workflow.Publication.IsPublished;
// Add the initially referenced definition
referencedWorkflowDefinitionsPublished[referencedDefinition.DefinitionId] = referencedDefinition;
// Build dependency map for topological sorting
var dependencyMap = filteredWorkflowReferences
.SelectMany(r => r.ReferencingDefinitionIds.Select(id => (id, r.ReferencedDefinitionId)))
.ToLookup(x => x.id, x => x.ReferencedDefinitionId);
// Perform topological sort to ensure dependent workflows are processed in the right order
var sortedWorkflowIds = referencingIds
.TSort(id => dependencyMap[id], true)
// Only process workflows that exist in our referencing workflows dictionary
.Where(id => referencingWorkflowGraphs.ContainsKey(id))
.ToList();
var updatedWorkflows = new Dictionary<string, UpdatedWorkflowDefinition>();
// Create a cache for drafts that we've already created during this operation
var draftCache = new Dictionary<string, WorkflowDefinition>();
foreach (var id in sortedWorkflowIds)
{
if (!referencingWorkflowGraphs.TryGetValue(id, out var graph) || !dependencyMap[id].Any())
continue;
foreach (var refId in dependencyMap[id])
{
var target = referencedWorkflowDefinitionsPublished.GetValueOrDefault(refId);
if (target == null) continue;
var updated = await UpdateWorkflowAsync(graph, target, draftCache, initialPublicationState, cancellationToken);
if (updated == null) continue;
graph = updated.NewGraph;
updatedWorkflows[updated.Definition.DefinitionId] = updated;
referencedWorkflowDefinitionsPublished[id] = updated.Definition;
draftCache[id] = updated.Definition;
referencingWorkflowGraphs[id] = updated.NewGraph;
}
}
_isUpdating = true;
foreach (var updatedWorkflow in updatedWorkflows.Values)
{
var requiresPublication = initialPublicationState.GetValueOrDefault(updatedWorkflow.Definition.DefinitionId);
if (requiresPublication)
await publisher.PublishAsync(updatedWorkflow.Definition, cancellationToken);
else
await publisher.SaveDraftAsync(updatedWorkflow.Definition, cancellationToken);
}
_isUpdating = false;
return new(updatedWorkflows.Select(u => u.Value.Definition));
}
private async IAsyncEnumerable<WorkflowReferences> GetReferencingWorkflowDefinitionIdsAsync(
string definitionId,
[EnumeratorCancellation] CancellationToken cancellationToken,
HashSet<string>? visitedIds = null)
{
visitedIds ??= new();
// If we've already processed this definition ID, skip it to prevent infinite recursion.
if (!visitedIds.Add(definitionId))
yield break;
var refs = (await workflowReferenceQuery.ExecuteAsync(definitionId, cancellationToken)).ToList();
yield return new(definitionId, refs);
foreach (var id in refs)
{
await foreach (var child in GetReferencingWorkflowDefinitionIdsAsync(id, cancellationToken, visitedIds))
yield return child;
}
}
private async Task<UpdatedWorkflowDefinition?> UpdateWorkflowAsync(
WorkflowGraph graph,
WorkflowDefinition target,
Dictionary<string, WorkflowDefinition> draftCache,
Dictionary<string, bool> initialPublicationState,
CancellationToken cancellationToken)
{
var willTargetBePublished = initialPublicationState.GetValueOrDefault(target.DefinitionId, target.IsPublished);
if (!willTargetBePublished)
return null;
var id = graph.Workflow.Identity.DefinitionId;
var draft = await GetOrCreateDraftAsync(id, draftCache, cancellationToken);
if (draft == null) return null;
var newGraph = await workflowDefinitionService.MaterializeWorkflowAsync(draft, cancellationToken);
var outdated = FindActivities(newGraph.Root, target.DefinitionId)
.Where(a => a.WorkflowDefinitionVersionId != target.Id)
.ToList();
if (!outdated.Any()) return null;
foreach (var act in outdated)
{
act.WorkflowDefinitionVersionId = target.Id;
act.Version = target.Version;
act.LatestAvailablePublishedVersionId = target.Id;
act.LatestAvailablePublishedVersion = target.Version;
}
if (newGraph.Root.Activity is Workflow wf)
draft.StringData = serializer.Serialize(wf.Root);
return new(draft, newGraph);
}
private async Task<WorkflowDefinition?> GetOrCreateDraftAsync(
string definitionId,
Dictionary<string, WorkflowDefinition> draftCache,
CancellationToken cancellationToken)
{
// Check if we already have a draft for this workflow
if (draftCache.TryGetValue(definitionId, out var cachedDraft))
return cachedDraft;
// Create or get a draft for this workflow
var draft = await publisher.GetDraftAsync(definitionId, VersionOptions.Latest, cancellationToken);
if (draft == null) return null;
// Store the draft in the cache for potential future use
draftCache[definitionId] = draft;
// Get the current published version of the workflow definition.
var publishedVersion = await workflowDefinitionStore.FindAsync(
WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Published).ToFilter(),
cancellationToken);
// Update the activity registry to be able to materialize the workflow.
var activityDescriptor = workflowDefinitionActivityDescriptorFactory.CreateDescriptor(draft, publishedVersion);
activityRegistry.Add(typeof(WorkflowDefinitionActivityProvider), activityDescriptor);
return draft;
}
private static IEnumerable<WorkflowDefinitionActivity> FindActivities(ActivityNode node, string definitionId)
{
// Do not drill into activities that are WorkflowDefinitionActivity
if (node.Activity is WorkflowDefinitionActivity)
yield break;
foreach (var child in node.Children)
{
if (child.Activity is WorkflowDefinitionActivity activity && activity.WorkflowDefinitionId == definitionId)
yield return activity;
foreach (var grandChildActivity in FindActivities(child, definitionId))
yield return grandChildActivity;
}
}
}

View file

@ -38,6 +38,12 @@ public class RunTask : Activity<object>
[Input(Description = "Any additional parameters to send to the task.")]
public Input<IDictionary<string, object>?> Payload { get; set; } = null!;
/// <summary>
/// The ID of the task that was requested to run.
/// </summary>
[Output(Description = "The ID of the task that was requested to run.")]
public Output<string> TaskId { get; set; } = null!;
/// <inheritdoc />
[JsonConstructor]
private RunTask(string? source = null, int? line = null) : base(source, line)
@ -95,6 +101,10 @@ public class RunTask : Activity<object>
var runTaskRequest = new RunTaskRequest(context, taskId, taskName, taskParams);
var dispatcher = context.GetRequiredService<ITaskDispatcher>();
// Set the task ID output.
TaskId.Set(context, taskId);
// Dispatch the task request.
await dispatcher.DispatchAsync(runTaskRequest, context.CancellationToken);
}