Update ForEach activity with local variables (#4315)

* Update ForEach activity with local variables

* Add option to include activity instance ID as part of bookmark hash
This commit is contained in:
Sipke Schoorstra 2023-08-11 21:59:39 +02:00 committed by GitHub
parent 57d06244ec
commit 3ba6a548e1
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
18 changed files with 230 additions and 154 deletions

View file

@ -43,23 +43,42 @@ public class DapperWorkflowInboxStore : IWorkflowInboxStore
var records = await _store.FindManyAsync(q => ApplyFilter(q, filter), cancellationToken);
return Map(records);
}
/// <inheritdoc />
public async ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(IEnumerable<WorkflowInboxMessageFilter> filters, CancellationToken cancellationToken = default)
{
var records = await _store.FindManyAsync(q => ApplyFilter(q, filters.ToArray()), cancellationToken);
return Map(records);
}
/// <inheritdoc />
public async ValueTask<long> DeleteAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default)
{
return await _store.DeleteAsync(q => ApplyFilter(q, filter), cancellationToken);
}
private void ApplyFilter(ParameterizedQuery query, WorkflowInboxMessageFilter filter)
private void ApplyFilter(ParameterizedQuery query, params WorkflowInboxMessageFilter[] filters)
{
query
.Is(nameof(WorkflowInboxMessageRecord.Hash), filter.Hash)
.Is(nameof(WorkflowInboxMessageRecord.WorkflowInstanceId), filter.WorkflowInstanceId)
.Is(nameof(WorkflowInboxMessageRecord.CorrelationId), filter.CorrelationId)
.Is(nameof(WorkflowInboxMessageRecord.ActivityTypeName), filter.ActivityTypeName)
.Is(nameof(WorkflowInboxMessageRecord.ActivityInstanceId), filter.ActivityInstanceId)
.Is(nameof(WorkflowInboxMessageRecord.IsHandled), filter.IsHandled)
;
var clauses = new List<ParameterizedQuery>();
foreach (var filter in filters)
{
var clause = new ParameterizedQuery(query.Dialect);
clause
.Is(nameof(WorkflowInboxMessageRecord.Hash), filter.Hash)
.Is(nameof(WorkflowInboxMessageRecord.WorkflowInstanceId), filter.WorkflowInstanceId)
.Is(nameof(WorkflowInboxMessageRecord.CorrelationId), filter.CorrelationId)
.Is(nameof(WorkflowInboxMessageRecord.ActivityTypeName), filter.ActivityTypeName)
.Is(nameof(WorkflowInboxMessageRecord.ActivityInstanceId), filter.ActivityInstanceId)
.Is(nameof(WorkflowInboxMessageRecord.IsHandled), filter.IsHandled)
;
clauses.Add(clause);
}
var clausesSql = string.Join(" OR ", $"({clauses.Select(x => x.Sql)})");
query.Sql.AppendLine(clausesSql);
}
private IEnumerable<WorkflowInboxMessage> Map(IEnumerable<WorkflowInboxMessageRecord> source) => source.Select(Map);

View file

@ -30,6 +30,16 @@ public class EFCoreWorkflowInboxStore : IWorkflowInboxStore
/// <inheritdoc />
public async ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default) => await _store.QueryAsync(filter.Apply, LoadAsync, cancellationToken);
/// <inheritdoc />
public async ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(IEnumerable<WorkflowInboxMessageFilter> filters, CancellationToken cancellationToken = default)
{
return await _store.QueryAsync(query =>
{
foreach (var filter in filters) filter.Apply(query);
return query;
}, LoadAsync, cancellationToken);
}
/// <inheritdoc />
public async ValueTask<long> DeleteAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default) => await _store.DeleteWhereAsync(filter.Apply, cancellationToken);

View file

@ -22,17 +22,28 @@ public class MongoWorkflowInboxStore : IWorkflowInboxStore
}
/// <inheritdoc />
public async ValueTask SaveAsync(WorkflowInboxMessage record, CancellationToken cancellationToken = default) =>
public async ValueTask SaveAsync(WorkflowInboxMessage record, CancellationToken cancellationToken = default) =>
await _mongoDbStore.SaveAsync(record, s => s.Id, cancellationToken);
/// <inheritdoc />
public async ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default) =>
(await _mongoDbStore.FindManyAsync(query => Filter(query, filter), cancellationToken));
public async ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default)
{
return await _mongoDbStore.FindManyAsync(query => Filter(query, filter), cancellationToken);
}
/// <inheritdoc />
public async ValueTask<long> DeleteAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default) =>
public async ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(IEnumerable<WorkflowInboxMessageFilter> filters, CancellationToken cancellationToken = default)
{
return await _mongoDbStore.FindManyAsync(query => Filter(query, filters.ToArray()), cancellationToken);
}
/// <inheritdoc />
public async ValueTask<long> DeleteAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default) =>
await _mongoDbStore.DeleteWhereAsync<string>(query => Filter(query, filter), x => x.Id, cancellationToken);
private IMongoQueryable<WorkflowInboxMessage> Filter(IMongoQueryable<WorkflowInboxMessage> queryable, WorkflowInboxMessageFilter filter) =>
(filter.Apply(queryable) as IMongoQueryable<WorkflowInboxMessage>)!;
private static IMongoQueryable<WorkflowInboxMessage> Filter(IMongoQueryable<WorkflowInboxMessage> queryable, params WorkflowInboxMessageFilter[] filters)
{
foreach (var filter in filters) filter.Apply(queryable);
return queryable;
}
}

View file

@ -1,10 +1,4 @@
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Extensions;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Behaviors;
using Elsa.Workflows.Core.Contracts;
using Elsa.Workflows.Core.Models;
using JetBrains.Annotations;
namespace Elsa.Workflows.Core.Activities;
@ -13,123 +7,6 @@ namespace Elsa.Workflows.Core.Activities;
/// Iterate over a set of values.
/// </summary>
[Activity("Elsa", "Looping", "Iterate over a set of values.")]
[PublicAPI]
public class ForEach : Activity
public class ForEach : ForEach<object>
{
private const string CurrentIndexProperty = "CurrentIndex";
/// <inheritdoc />
public ForEach([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
{
Behaviors.Add<BreakBehavior>(this);
}
/// <inheritdoc />
public ForEach(ICollection<object> items, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Items = new Input<ICollection<object>>(items);
}
/// <summary>
/// The set of values to iterate.
/// </summary>
[Input(Description = "The set of values to iterate.")]
public Input<ICollection<object>> Items { get; set; } = new(Array.Empty<object>());
/// <summary>
/// The activity to execute for each iteration.
/// </summary>
[Port]
public IActivity? Body { get; set; }
/// <summary>
/// The current value being iterated will be assigned to the specified <see cref="MemoryBlockReference"/>.
/// </summary>
[Output(Description = "Assign the current value to the specified variable.")]
public Output<object>? CurrentValue { get; set; }
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
// Execute first iteration.
await HandleIteration(context);
}
private async Task HandleIteration(ActivityExecutionContext context)
{
var currentIndex = context.GetProperty<int>(CurrentIndexProperty);
var items = context.Get(Items)!.ToList();
if (currentIndex >= items.Count)
{
await context.CompleteActivityAsync();
return;
}
var currentItem = items[currentIndex];
context.Set(CurrentValue, currentItem);
if (Body != null)
await context.ScheduleActivityAsync(Body, OnChildCompleted);
else
await context.CompleteActivityAsync();
// Increment index.
context.UpdateProperty<int>(CurrentIndexProperty, x => x + 1);
}
private async ValueTask OnChildCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext)
{
await HandleIteration(context);
}
}
/// <summary>
/// A strongly-typed for-each construct where <see cref="T"/> is the item type.
/// </summary>
[PublicAPI]
public class ForEach<T> : ForEach
{
/// <inheritdoc />
public ForEach([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <inheritdoc />
public ForEach(Input<ICollection<T>> items) : this()
{
Items = items;
}
/// <inheritdoc />
public ForEach(ICollection<T> items) : this(new Input<ICollection<T>>(items))
{
}
/// <inheritdoc />
public ForEach(Func<ExpressionExecutionContext, ICollection<T>> items) : this(new Input<ICollection<T>>(items))
{
}
/// <summary>
/// The items to iterate on.
/// </summary>
[Input]
public new Input<ICollection<T>> Items
{
get => new(base.Items.Expression, base.Items.MemoryBlockReference());
set => base.Items = new Input<ICollection<object>>(value.Expression, value.MemoryBlockReference());
}
/// <summary>
/// Provides access to the current value.
/// </summary>
public new Output<T?> CurrentValue
{
get =>
base.CurrentValue != null
? new(base.CurrentValue.MemoryBlockReference)
: new();
set => base.CurrentValue = new Output<object>(value.MemoryBlockReference);
}
}

View file

@ -0,0 +1,105 @@
using System.Runtime.CompilerServices;
using Elsa.Expressions.Models;
using Elsa.Extensions;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Behaviors;
using Elsa.Workflows.Core.Contracts;
using Elsa.Workflows.Core.Memory;
using Elsa.Workflows.Core.Models;
namespace Elsa.Workflows.Core.Activities;
/// <summary>
/// A strongly-typed for-each construct where <see cref="T"/> is the item type.
/// </summary>
public class ForEach<T> : Activity
{
private const string CurrentIndexProperty = "CurrentIndex";
/// <inheritdoc />
public ForEach(Func<ExpressionExecutionContext, ICollection<T>> @delegate, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input<ICollection<T>>(@delegate), source, line)
{
}
/// <inheritdoc />
public ForEach(Func<ICollection<T>> @delegate, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input<ICollection<T>>(@delegate), source, line)
{
}
/// <inheritdoc />
public ForEach(ICollection<T> items, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input<ICollection<T>>(items), source, line)
{
}
/// <inheritdoc />
public ForEach(Input<ICollection<T>> items, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line)
{
Items = items;
}
/// <inheritdoc />
public ForEach([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default)
{
Behaviors.Add<BreakBehavior>(this);
}
/// <summary>
/// The set of values to iterate.
/// </summary>
[Input(Description = "The set of values to iterate.")]
public Input<ICollection<T>> Items { get; set; } = new(Array.Empty<T>());
/// <summary>
/// The activity to execute for each iteration.
/// </summary>
[Port]
public IActivity? Body { get; set; }
/// <summary>
/// The current value being iterated will be assigned to the specified <see cref="MemoryBlockReference"/>.
/// </summary>
[Output(Description = "Assign the current value to the specified variable.")]
public Output<T>? CurrentValue { get; set; }
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
// Execute first iteration.
await HandleIteration(context);
}
private async Task HandleIteration(ActivityExecutionContext context)
{
var currentIndex = context.GetProperty<int>(CurrentIndexProperty);
var items = context.Get(Items)!.ToList();
if (currentIndex >= items.Count)
{
await context.CompleteActivityAsync();
return;
}
var currentValue = items[currentIndex];
context.Set(CurrentValue, currentValue);
if (Body != null)
{
var variables = new[]
{
new Variable("CurrentIndex", currentIndex),
new Variable("CurrentValue", currentValue)
};
await context.ScheduleActivityAsync(Body, OnChildCompleted, variables: variables);
}
else
await context.CompleteActivityAsync();
// Increment index.
context.UpdateProperty<int>(CurrentIndexProperty, x => x + 1);
}
private async ValueTask OnChildCompleted(ActivityExecutionContext context, ActivityExecutionContext childContext)
{
await HandleIteration(context);
}
}

View file

@ -45,7 +45,7 @@ public class ParallelForEach<T> : Activity
foreach (var item in items)
{
// For each item, declare a new variable the work to be scheduled.
var variable = new Variable<T>("CurrentItem", item)
var variable = new Variable<T>("CurrentValue", item)
{
// TODO: This should be configurable, because this won't work for e.g. file streams and other non-serializable types.
StorageDriverType = typeof(WorkflowStorageDriver)

View file

@ -314,7 +314,8 @@ public class ActivityExecutionContext : IExecutionContext
var bookmarkName = options?.BookmarkName ?? Activity.Type;
var bookmarkHasher = GetRequiredService<IBookmarkHasher>();
var identityGenerator = GetRequiredService<IIdentityGenerator>();
var hash = bookmarkHasher.Hash(bookmarkName, payload);
var includeActivityInstanceId = options?.IncludeActivityInstanceId ?? true;
var hash = bookmarkHasher.Hash(bookmarkName, payload, includeActivityInstanceId ? Id : null);
var bookmark = new Bookmark(
identityGenerator.GenerateId(),

View file

@ -6,7 +6,7 @@ namespace Elsa.Workflows.Core.Contracts;
public interface IBookmarkHasher
{
/// <summary>
/// Produces a hash from the specified activity type name and bookmark payload.
/// Produces a hash from the specified activity type name, bookmark payload and activity instance ID.
/// </summary>
string Hash(string activityTypeName, object? payload);
string Hash(string activityTypeName, object? payload, string? activityInstanceId = default);
}

View file

@ -71,7 +71,7 @@ public static class WorkflowExecutionContextExtensions
/// </summary>
public static void ScheduleBookmark(this WorkflowExecutionContext workflowExecutionContext, Bookmark bookmark)
{
// Construct bookmark.
// Get the activity execution context that owns the bookmark.
var bookmarkedActivityContext = workflowExecutionContext.ActiveActivityExecutionContexts.FirstOrDefault(x => x.Id == bookmark.ActivityInstanceId);
if(bookmarkedActivityContext == null)

View file

@ -2,4 +2,12 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Models;
public record BookmarkOptions(object? Payload = default, ExecuteActivityDelegate? Callback = default, string? BookmarkName = default, bool AutoBurn = true);
/// <summary>
/// Provides bookmark creation options.
/// </summary>
/// <param name="Payload">An optional payload to associate with the bookmark.</param>
/// <param name="Callback">An optional callback to invoke when the bookmark is triggered.</param>
/// <param name="BookmarkName">An optional name to associate with the bookmark.</param>
/// <param name="AutoBurn">Whether or not the bookmark should be automatically burned when triggered.</param>
/// <param name="IncludeActivityInstanceId">Whether or not the activity instance ID should be included in the bookmark payload.</param>
public record BookmarkOptions(object? Payload = default, ExecuteActivityDelegate? Callback = default, string? BookmarkName = default, bool AutoBurn = true, bool IncludeActivityInstanceId = true);

View file

@ -1,3 +1,4 @@
using System.Text;
using System.Text.Json;
using Elsa.Expressions.Contracts;
using Elsa.Workflows.Core.Contracts;
@ -30,10 +31,18 @@ public class BookmarkHasher : IBookmarkHasher
}
/// <inheritdoc />
public string Hash(string activityTypeName, object? payload)
public string Hash(string activityTypeName, object? payload, string? activityInstanceId = default)
{
var json = payload != null ? Serialize(payload) : null;
var input = $"{activityTypeName}{(!string.IsNullOrWhiteSpace(json) ? ":" + json : "")}";
var inputSource = new List<string> { activityTypeName};
if (!string.IsNullOrWhiteSpace(json))
inputSource.Add(json);
if (!string.IsNullOrWhiteSpace(activityInstanceId))
inputSource.Add(activityInstanceId);
var input = string.Join("|", inputSource);
var hash = _hasher.Hash(input);
return hash;

View file

@ -75,7 +75,12 @@ public class Event : Trigger<object?>
if (!context.IsTriggerOfWorkflow())
{
context.CreateBookmark(new EventBookmarkPayload(eventName));
var options = new BookmarkOptions
{
Payload = new EventBookmarkPayload(eventName),
IncludeActivityInstanceId = false
};
context.CreateBookmark(options);
return;
}

View file

@ -40,4 +40,12 @@ public interface IWorkflowInbox
/// <param name="cancellationToken">An optional cancellation token.</param>
/// <returns>A list of messages that match the filter.</returns>
ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default);
/// <summary>
/// Finds all messages matching the specified filter.
/// </summary>
/// <param name="filters">The filters to apply.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
/// <returns>A list of messages that match the filter.</returns>
ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(IEnumerable<WorkflowInboxMessageFilter> filters, CancellationToken cancellationToken = default);
}

View file

@ -24,6 +24,14 @@ public interface IWorkflowInboxStore
/// <param name="cancellationToken">An optional cancellation token.</param>
/// <returns>A list of messages that match the filter.</returns>
ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default);
/// <summary>
/// Finds all messages matching the specified filters.
/// </summary>
/// <param name="filters">The filters to apply.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
/// <returns>A list of messages that match the filter.</returns>
ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(IEnumerable<WorkflowInboxMessageFilter> filters, CancellationToken cancellationToken = default);
/// <summary>
/// Deletes all messages matching the specified filter.

View file

@ -28,12 +28,10 @@ public class DeliverWorkflowMessagesFromInbox : INotificationHandler<WorkflowBoo
foreach (var bookmark in addedBookmarks)
{
var activityTypeName = bookmark.Name;
var activityInstanceId = bookmark.ActivityInstanceId;
var hash = bookmark.Hash;
var filter = new WorkflowInboxMessageFilter
{
ActivityTypeName = activityTypeName,
ActivityInstanceId = activityInstanceId,
Hash = hash,
IsHandled = false
};

View file

@ -109,4 +109,10 @@ public class DefaultWorkflowInbox : IWorkflowInbox
{
return await _store.FindManyAsync(filter, cancellationToken);
}
/// <inheritdoc />
public async ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(IEnumerable<WorkflowInboxMessageFilter> filters, CancellationToken cancellationToken = default)
{
return await _store.FindManyAsync(filters, cancellationToken);
}
}

View file

@ -34,6 +34,13 @@ public class MemoryWorkflowInboxStore : IWorkflowInboxStore
return new(entities);
}
/// <inheritdoc />
public ValueTask<IEnumerable<WorkflowInboxMessage>> FindManyAsync(IEnumerable<WorkflowInboxMessageFilter> filters, CancellationToken cancellationToken = default)
{
var entities = _store.Query(query => Filter(query, filters.ToArray())).ToList();
return new(entities);
}
/// <inheritdoc />
public async ValueTask<long> DeleteAsync(WorkflowInboxMessageFilter filter, CancellationToken cancellationToken = default)
{
@ -41,5 +48,9 @@ public class MemoryWorkflowInboxStore : IWorkflowInboxStore
return _store.DeleteMany(ids);
}
private static IQueryable<WorkflowInboxMessage> Filter(IQueryable<WorkflowInboxMessage> query, WorkflowInboxMessageFilter filter) => filter.Apply(query);
private static IQueryable<WorkflowInboxMessage> Filter(IQueryable<WorkflowInboxMessage> query, params WorkflowInboxMessageFilter[] filters)
{
foreach (var filter in filters) filter.Apply(query);
return query;
}
}

View file

@ -19,9 +19,9 @@ public static class ForEachWorkflow
Activities =
{
new WriteLine("Going through the shopping list..."),
new ForEach(shoppingList)
new ForEach<string>(shoppingList)
{
CurrentValue = new Output<object>(currentValueVariable),
CurrentValue = new Output<string>(currentValueVariable),
Body = new Sequence
{
Activities =