diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkBoundWorkflowService.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkBoundWorkflowService.cs index 72c3e1f7d..72891723b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkBoundWorkflowService.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkBoundWorkflowService.cs @@ -5,6 +5,7 @@ namespace Elsa.Workflows.Runtime; /// /// Represents a service that looks up bookmark-bound workflows. /// +[Obsolete("Will be removed in a future version.")] public interface IBookmarkBoundWorkflowService { /// diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs index 865553096..8e62b5688 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkResumer.cs @@ -6,6 +6,7 @@ namespace Elsa.Workflows.Runtime; /// /// Resumes workflows using a given stimulus or bookmark filter. /// +[Obsolete("Use IWorkflowResumer instead.")] public interface IBookmarkResumer { /// diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowResumer.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowResumer.cs new file mode 100644 index 000000000..ce1b56e06 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowResumer.cs @@ -0,0 +1,39 @@ +using Elsa.Workflows.Runtime.Filters; +using Elsa.Workflows.Runtime.Messages; +using Elsa.Workflows.Runtime.Options; + +namespace Elsa.Workflows.Runtime; + +/// +/// Resumes workflows using a given stimulus or bookmark filter. +/// +public interface IWorkflowResumer +{ + /// + /// Resumes the workflows associated with the bookmarks matching the given stimulus. + /// + Task> ResumeAsync(object stimulus, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity; + + /// + /// Resumes the workflow associated with the bookmark specified by the given bookmark ID. + /// + Task ResumeAsync(string bookmarkId, IDictionary input, CancellationToken cancellationToken = default); + + /// + /// Resumes the workflows associated with the bookmarks matching the given stimulus. If a workflow instance ID is specified, only resumes workflows associated with that instance. + /// + Task> ResumeAsync(object stimulus, string? workflowInstanceId, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity; + + /// + /// Resumes the workflow associated with the bookmark specified by the given bookmark ID. + /// + Task ResumeAsync(string bookmarkId, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity; + + /// Resumes the workflows associated with the bookmarks matching the given request. + Task> ResumeAsync(ResumeBookmarkRequest request, CancellationToken cancellationToken = default); + + /// + /// Resumes the workflows matching the given bookmark filter. + /// + Task> ResumeAsync(BookmarkFilter filter, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index d8f57b739..ccde9984f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -276,6 +276,7 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module) .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkFilter.cs b/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkFilter.cs index 544c5ddb7..b632362f3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkFilter.cs +++ b/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkFilter.cs @@ -1,3 +1,5 @@ +using System.Collections; +using System.Text; using Elsa.Workflows.Runtime.Entities; namespace Elsa.Workflows.Runtime.Filters; @@ -7,6 +9,9 @@ namespace Elsa.Workflows.Runtime.Filters; /// public class BookmarkFilter { + // Cache the properties of BookmarkFilter for performance. + private static readonly System.Reflection.PropertyInfo[] CachedProperties = typeof(BookmarkFilter).GetProperties(); + /// /// Gets or sets the ID of the bookmark. /// @@ -86,4 +91,40 @@ public class BookmarkFilter { Names = activityTypeNames.ToList() }; + + public string GetHashableString() + { + // Return a hashable string representation of the filter, excluding null values. + var sb = new StringBuilder(); + foreach (var prop in CachedProperties) + { + var value = prop.GetValue(this); + if (value == null) + continue; + + string valueString; + // Handle collections (excluding string) + if (value is IEnumerable enumerable and not string) + { + var items = new List(); + foreach (var item in enumerable) + { + if (item != null) + items.Add(item.ToString()!); + } + items.Sort(StringComparer.Ordinal); + valueString = string.Join(",", items); + } + else + { + var toStringResult = value.ToString(); + if (toStringResult == null) + continue; + valueString = toStringResult; + } + sb.Append($"{prop.Name}:{valueString};"); + } + + return sb.ToString(); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs b/src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs index 58a29bed6..4688ae0bb 100644 --- a/src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs +++ b/src/modules/Elsa.Workflows.Runtime/Requests/ResumeBookmarkRequest.cs @@ -4,14 +4,20 @@ namespace Elsa.Workflows.Runtime; public class ResumeBookmarkRequest { - public string WorkflowInstanceId { get; set; } = default!; + public string WorkflowInstanceId { get; set; } = null!; /// The ID of the bookmark that triggered the workflow instance, if any. - public string BookmarkId { get; set; } = default!; + public string BookmarkId { get; set; } = null!; /// The handle of the activity to schedule, if any. + [Obsolete("Use ActivityInstanceId instead")] public ActivityHandle? ActivityHandle { get; set; } + /// + /// The ID of the activity instance to resume, if any. + /// + public string? ActivityInstanceId { get; set; } + /// Any additional properties to associate with the workflow instance. public IDictionary? Properties { get; set; } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkBoundWorkflowService.cs b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkBoundWorkflowService.cs index 7fc4e2c4b..9c4c58561 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkBoundWorkflowService.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkBoundWorkflowService.cs @@ -4,6 +4,7 @@ using Elsa.Workflows.Runtime.Options; namespace Elsa.Workflows.Runtime; /// +[Obsolete("Will be removed in a future version.")] public class BookmarkBoundWorkflowService(IWorkflowMatcher workflowMatcher) : IBookmarkBoundWorkflowService { /// diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkQueueProcessor.cs b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkQueueProcessor.cs index 425b4c63f..341e0ef69 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkQueueProcessor.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkQueueProcessor.cs @@ -7,7 +7,7 @@ using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Runtime; -public class BookmarkQueueProcessor(IBookmarkQueueStore store, IBookmarkResumer bookmarkResumer, ILogger logger) : IBookmarkQueueProcessor +public class BookmarkQueueProcessor(IBookmarkQueueStore store, IWorkflowResumer workflowResumer, ILogger logger) : IBookmarkQueueProcessor { public async Task ProcessAsync(CancellationToken cancellationToken = default) { @@ -41,16 +41,16 @@ public class BookmarkQueueProcessor(IBookmarkQueueStore store, IBookmarkResumer logger.LogDebug("Processing bookmark queue item {BookmarkQueueItemId} for workflow instance {WorkflowInstanceId} for activity type {ActivityType}", item.Id, item.WorkflowInstanceId, item.ActivityTypeName); - var result = await bookmarkResumer.ResumeAsync(filter, options, cancellationToken); + var responses = (await workflowResumer.ResumeAsync(filter, options, cancellationToken)).ToList(); - if (result.Matched) + if (responses.Count > 0) { - logger.LogDebug("Successfully resumed workflow instance {WorkflowInstance} using bookmark {BookmarkId} for activity type {ActivityType}", item.WorkflowInstanceId, item.BookmarkId, item.ActivityTypeName); + logger.LogDebug("Successfully resumed {WorkflowCount} workflow instances using stimulus {StimulusHash} for activity type {ActivityType}", responses.Count, item.StimulusHash, item.ActivityTypeName); await store.DeleteAsync(item.Id, cancellationToken); } else { - logger.LogDebug("No matching bookmark found for bookmark queue item {BookmarkQueueItemId} for workflow instance {WorkflowInstanceId} for activity type {ActivityType}", item.Id, item.WorkflowInstanceId, item.ActivityTypeName); + logger.LogDebug("No matching bookmarks found for bookmark queue item {BookmarkQueueItemId} for workflow instance {WorkflowInstanceId} for activity type {ActivityType} with stimulus {StimulusHash}", item.Id, item.WorkflowInstanceId, item.ActivityTypeName, item.StimulusHash); } } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs index d91c37aa0..206cd7244 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkResumer.cs @@ -8,6 +8,7 @@ using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Runtime; /// +[Obsolete("Use WorkflowResumer instead.")] public class BookmarkResumer(IWorkflowRuntime workflowRuntime, IBookmarkStore bookmarkStore, IStimulusHasher stimulusHasher, ILogger logger) : IBookmarkResumer { /// diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs b/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs index bbb688550..2d93bb75b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs @@ -1,6 +1,5 @@ -using Elsa.Workflows.Models; +using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.Messages; -using Elsa.Workflows.Runtime.Options; using Elsa.Workflows.Runtime.Results; using Microsoft.Extensions.Logging; using Open.Linq.AsyncExtensions; @@ -11,9 +10,8 @@ namespace Elsa.Workflows.Runtime; public class StimulusSender( IStimulusHasher stimulusHasher, ITriggerBoundWorkflowService triggerBoundWorkflowService, - IBookmarkBoundWorkflowService bookmarkBoundWorkflowService, + IWorkflowResumer workflowResumer, IBookmarkQueue bookmarkQueue, - IWorkflowRuntime workflowRuntime, ITriggerInvoker triggerInvoker, ILogger logger) : IStimulusSender { @@ -65,15 +63,15 @@ public class StimulusSender( Properties = properties, ParentWorkflowInstanceId = parentId }; - + var response = await triggerInvoker.InvokeAsync(triggerRequest, cancellationToken); - + if (response.CannotStart) { logger.LogWarning("Workflow activation strategy disallowed starting workflow {WorkflowDefinitionHandle} with correlation ID {CorrelationId}", workflow.DefinitionHandle, correlationId); continue; } - + responses.Add(response.ToRunWorkflowInstanceResponse()); } } @@ -83,60 +81,48 @@ public class StimulusSender( private async Task> ResumeExistingWorkflowsAsync(string stimulusHash, StimulusMetadata? metadata, CancellationToken cancellationToken) { - var bookmarkOptions = metadata != null - ? new FindBookmarkOptions - { - CorrelationId = metadata.CorrelationId, - WorkflowInstanceId = metadata.WorkflowInstanceId, - ActivityInstanceId = metadata.ActivityInstanceId, - } - : null; - var bookmarkBoundWorkflows = await bookmarkBoundWorkflowService.FindManyAsync(stimulusHash, bookmarkOptions, cancellationToken).ToList(); var input = metadata?.Input; var properties = metadata?.Properties; - var activityHandle = metadata?.ActivityInstanceId != null ? ActivityHandle.FromActivityInstanceId(metadata.ActivityInstanceId) : null; - var responses = new List(); - - if (bookmarkBoundWorkflows.Count > 0) + + var bookmarkFilter = new BookmarkFilter { - foreach (var bookmarkBoundWorkflow in bookmarkBoundWorkflows) - { - var workflowInstanceId = bookmarkBoundWorkflow.WorkflowInstanceId; - var workflowClient = await workflowRuntime.CreateClientAsync(workflowInstanceId, cancellationToken); + Hash = stimulusHash, + CorrelationId = metadata?.CorrelationId, + WorkflowInstanceId = metadata?.WorkflowInstanceId, + ActivityInstanceId = metadata?.ActivityInstanceId, + BookmarkId = metadata?.BookmarkId + }; + var responses = (await workflowResumer.ResumeAsync(bookmarkFilter, new() + { + Input = input, + Properties = properties + }, cancellationToken)).ToList(); - foreach (var storedBookmark in bookmarkBoundWorkflow.Bookmarks) - { - var request = new RunWorkflowInstanceRequest - { - Input = input, - Properties = properties, - ActivityHandle = activityHandle, - BookmarkId = storedBookmark.Id, - }; - var response = await workflowClient.RunInstanceAsync(request, cancellationToken); - responses.Add(response); - } + if (responses.Count > 0) + { + logger.LogDebug("Successfully resumed {WorkflowCount} workflow instances using stimulus {StimulusHash}", responses.Count, stimulusHash); + return responses; + } + + // If no bookmarks were matched, enqueue the request in case a matching bookmark is created in the near future. + var workflowInstanceId = metadata?.WorkflowInstanceId; + + var bookmarkQueueItem = new NewBookmarkQueueItem + { + WorkflowInstanceId = workflowInstanceId, + BookmarkId = metadata?.BookmarkId, + CorrelationId = metadata?.CorrelationId, + StimulusHash = stimulusHash, + Options = new() + { + Input = input, + Properties = properties } - } - else - { - // If no bookmarks were matched, enqueue the request in case a matching bookmark is created in the near future. - var workflowInstanceId = metadata?.WorkflowInstanceId; - - var bookmarkQueueItem = new NewBookmarkQueueItem - { - WorkflowInstanceId = workflowInstanceId, - BookmarkId = metadata?.BookmarkId, - CorrelationId = metadata?.CorrelationId, - StimulusHash = stimulusHash, - Options = new() - { - Input = input, - Properties = properties - } - }; - await bookmarkQueue.EnqueueAsync(bookmarkQueueItem, cancellationToken); - } + }; + + logger.LogDebug("Bookmark queue item enqueued with stimulus: {StimulusHash}", bookmarkQueueItem.StimulusHash); + + await bookmarkQueue.EnqueueAsync(bookmarkQueueItem, cancellationToken); return responses; } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreBookmarkQueue.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreBookmarkQueue.cs index b3edde3f1..8a7c87894 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StoreBookmarkQueue.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreBookmarkQueue.cs @@ -1,13 +1,11 @@ using Elsa.Common; using Elsa.Workflows.Runtime.Entities; -using Elsa.Workflows.Runtime.Filters; using Microsoft.Extensions.Logging; namespace Elsa.Workflows.Runtime; public class StoreBookmarkQueue( IBookmarkQueueStore store, - IBookmarkResumer resumer, IBookmarkQueueSignaler bookmarkQueueSignaler, ISystemClock systemClock, IIdentityGenerator identityGenerator, @@ -15,26 +13,6 @@ public class StoreBookmarkQueue( { public async Task EnqueueAsync(NewBookmarkQueueItem item, CancellationToken cancellationToken = default) { - var filter = new BookmarkFilter - { - BookmarkId = item.BookmarkId, - CorrelationId = item.CorrelationId, - Hash = item.StimulusHash, - WorkflowInstanceId = item.WorkflowInstanceId, - Name = item.ActivityTypeName - }; - - var result = await resumer.ResumeAsync(filter, item.Options, cancellationToken); - - if (result.Matched) - { - logger.LogDebug("Successfully resumed workflow instance {WorkflowInstance} using bookmark {BookmarkId} for activity type {ActivityType}", item.WorkflowInstanceId, item.BookmarkId, item.ActivityTypeName); - return; - } - - // There was no matching bookmark yet, or the associated workflow instance hasn't been stored in the DB yet. Store the queue item for the system to pick up whenever the bookmark or workflow instance becomes present. - logger.LogDebug("No bookmark with ID {BookmarkId} found for workflow {WorkflowInstance} for activity type {ActivityType}. Adding the request to the bookmark queue", item.BookmarkId, item.WorkflowInstanceId, item.ActivityTypeName); - var entity = new BookmarkQueueItem { Id = identityGenerator.GenerateId(), @@ -48,6 +26,8 @@ public class StoreBookmarkQueue( CreatedAt = systemClock.UtcNow, }; + logger.LogDebug("Enqueuing bookmark queue item {BookmarkQueueItemId} with bookmark {BookmarkId} and stimulus {StimulusHash}", entity.Id, entity.BookmarkId, entity.StimulusHash); + await store.AddAsync(entity, cancellationToken); // Trigger the bookmark queue processor. diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowResumer.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowResumer.cs new file mode 100644 index 000000000..08b651227 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowResumer.cs @@ -0,0 +1,135 @@ +using Elsa.Common.DistributedHosting; +using Elsa.Workflows.Helpers; +using Elsa.Workflows.Runtime.Exceptions; +using Elsa.Workflows.Runtime.Filters; +using Elsa.Workflows.Runtime.Messages; +using Elsa.Workflows.Runtime.Options; +using Medallion.Threading; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; + +namespace Elsa.Workflows.Runtime; + +/// +public class WorkflowResumer( + IWorkflowRuntime workflowRuntime, + IBookmarkStore bookmarkStore, + IStimulusHasher stimulusHasher, + IDistributedLockProvider distributedLockProvider, + IOptions distributedLockingOptions, + ILogger logger) : IWorkflowResumer +{ + /// + public Task> ResumeAsync(object stimulus, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity + { + return ResumeAsync(stimulus, null, options, cancellationToken); + } + + /// + public async Task> ResumeAsync(object stimulus, string? workflowInstanceId = null, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity + { + var activityTypeName = ActivityTypeNameHelper.GenerateTypeName(); + var stimulusHash = stimulusHasher.Hash(activityTypeName, stimulus); + var bookmarkFilter = new BookmarkFilter + { + Name = activityTypeName, + WorkflowInstanceId = workflowInstanceId, + Hash = stimulusHash, + }; + return await ResumeAsync(bookmarkFilter, options, cancellationToken); + } + + /// + public async Task ResumeAsync(string bookmarkId, IDictionary input, CancellationToken cancellationToken = default) + { + var bookmarkFilter = new BookmarkFilter + { + BookmarkId = bookmarkId + }; + var options = new ResumeBookmarkOptions + { + Input = input + }; + var responses = await ResumeAsync(bookmarkFilter, options, cancellationToken); + return responses.FirstOrDefault(); + } + + /// + public async Task ResumeAsync(string bookmarkId, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity + { + var activityTypeName = ActivityTypeNameHelper.GenerateTypeName(); + var bookmarkFilter = new BookmarkFilter + { + Name = activityTypeName, + BookmarkId = bookmarkId + }; + var response = await ResumeAsync(bookmarkFilter, options, cancellationToken); + return response.FirstOrDefault(); + } + + public async Task> ResumeAsync(ResumeBookmarkRequest request, CancellationToken cancellationToken = default) + { + var filter = new BookmarkFilter + { + BookmarkId = request.BookmarkId, + ActivityInstanceId = request.ActivityInstanceId ?? request.ActivityHandle?.ActivityInstanceId, + }; + + var resumeOptions = new ResumeBookmarkOptions() + { + Input = request.Input, + Properties = request.Properties, + }; + return await ResumeAsync(filter, resumeOptions, cancellationToken); + } + + /// + public async Task> ResumeAsync(BookmarkFilter filter, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) + { + var hashableFilterString = filter.GetHashableString(); + var lockKey = $"workflow-resumer:{hashableFilterString}"; + + try + { + await using var filterLock = await distributedLockProvider.AcquireLockAsync(lockKey, distributedLockingOptions.Value.LockAcquisitionTimeout, cancellationToken); + var bookmarks = (await bookmarkStore.FindManyAsync(filter, cancellationToken)).ToList(); + + if (bookmarks.Count == 0) + { + logger.LogDebug("No bookmarks found in store for filter {@Filter}", filter); + return []; + } + + var responses = new List(); + foreach (var bookmark in bookmarks) + { + var workflowClient = await workflowRuntime.CreateClientAsync(bookmark.WorkflowInstanceId, cancellationToken); + var runRequest = new RunWorkflowInstanceRequest + { + Input = options?.Input, + Properties = options?.Properties, + BookmarkId = bookmark.Id + }; + + try + { + var response = await workflowClient.RunInstanceAsync(runRequest, cancellationToken); + logger.LogDebug("Resumed workflow instance {WorkflowInstanceId} with bookmark {BookmarkId}", bookmark.WorkflowInstanceId, bookmark.Id); + responses.Add(response); + } + catch (WorkflowInstanceNotFoundException) + { + // The workflow instance does not (yet) exist in the DB. + logger.LogDebug("No workflow instance with ID {WorkflowInstanceId} found for bookmark {BookmarkId} at this time.", bookmark.WorkflowInstanceId, bookmark.Id); + } + } + + return responses; + } + catch (TimeoutException e) + { + // Rethrow but with a more specific message. + throw new TimeoutException($"Could not acquire distributed lock with key '{lockKey}' within the configured timeout of {distributedLockingOptions.Value.LockAcquisitionTimeout}.", e); + } + } +} \ No newline at end of file