Fix race condition when sending same stimuli (#6895)
* Introduce `WorkflowResumer` service and deprecate `BookmarkResumer`. - Adds `IWorkflowResumer` and its implementation for workflow resumption. - Marks `BookmarkResumer` and related interfaces as obsolete. - Refactors dependent services to use `WorkflowResumer`. - Enhances `ResumeBookmarkRequest` to include `ActivityInstanceId`. - Updates logging and queue handling logic to align with the new resumption approach. * Update lock key prefix in `WorkflowResumer` for consistency with service naming. * Add exception handling for distributed lock acquisition in `WorkflowResumer` - Wrap distributed lock logic with `try-catch` to handle `TimeoutException`. - Improve error message when lock acquisition fails due to timeout. - Preserve existing workflow resumption behavior and logging. * Update src/modules/Elsa.Workflows.Runtime/Filters/BookmarkFilter.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Optimize `BookmarkFilter` hashing logic for improved performance and readability. * Merge remote-tracking branch 'origin/enh/locked-bookmark-resumption-2' into enh/locked-bookmark-resumption-2 * Remove unused variable and redundant line breaks for cleaner code. * Clean up logging configuration by removing unused debug log levels. * Update src/modules/Elsa.Workflows.Runtime/Services/WorkflowResumer.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Handle collections in `BookmarkFilter` hashing to ensure determinism and improve compatibility. * Refactor `BookmarkFilter` hashing logic for clarity and consistency. * Improve `TimeoutException` handling with a more descriptive message in `WorkflowResumer`. * Update src/modules/Elsa.Workflows.Runtime/Filters/BookmarkFilter.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Simplify `BookmarkFilter` by utilizing `using` directives and refining type references. --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
parent
6c714fe793
commit
c549f49dfb
|
|
@ -5,6 +5,7 @@ namespace Elsa.Workflows.Runtime;
|
|||
/// <summary>
|
||||
/// Represents a service that looks up bookmark-bound workflows.
|
||||
/// </summary>
|
||||
[Obsolete("Will be removed in a future version.")]
|
||||
public interface IBookmarkBoundWorkflowService
|
||||
{
|
||||
/// <summary>
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ namespace Elsa.Workflows.Runtime;
|
|||
/// <summary>
|
||||
/// Resumes workflows using a given stimulus or bookmark filter.
|
||||
/// </summary>
|
||||
[Obsolete("Use IWorkflowResumer instead.")]
|
||||
public interface IBookmarkResumer
|
||||
{
|
||||
/// <summary>
|
||||
|
|
|
|||
|
|
@ -0,0 +1,39 @@
|
|||
using Elsa.Workflows.Runtime.Filters;
|
||||
using Elsa.Workflows.Runtime.Messages;
|
||||
using Elsa.Workflows.Runtime.Options;
|
||||
|
||||
namespace Elsa.Workflows.Runtime;
|
||||
|
||||
/// <summary>
|
||||
/// Resumes workflows using a given stimulus or bookmark filter.
|
||||
/// </summary>
|
||||
public interface IWorkflowResumer
|
||||
{
|
||||
/// <summary>
|
||||
/// Resumes the workflows associated with the bookmarks matching the given stimulus.
|
||||
/// </summary>
|
||||
Task<IEnumerable<RunWorkflowInstanceResponse>> ResumeAsync<TActivity>(object stimulus, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity;
|
||||
|
||||
/// <summary>
|
||||
/// Resumes the workflow associated with the bookmark specified by the given bookmark ID.
|
||||
/// </summary>
|
||||
Task<RunWorkflowInstanceResponse?> ResumeAsync(string bookmarkId, IDictionary<string, object> input, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// 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.
|
||||
/// </summary>
|
||||
Task<IEnumerable<RunWorkflowInstanceResponse>> ResumeAsync<TActivity>(object stimulus, string? workflowInstanceId, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity;
|
||||
|
||||
/// <summary>
|
||||
/// Resumes the workflow associated with the bookmark specified by the given bookmark ID.
|
||||
/// </summary>
|
||||
Task<RunWorkflowInstanceResponse?> ResumeAsync<TActivity>(string bookmarkId, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity;
|
||||
|
||||
/// Resumes the workflows associated with the bookmarks matching the given request.
|
||||
Task<IEnumerable<RunWorkflowInstanceResponse>> ResumeAsync(ResumeBookmarkRequest request, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Resumes the workflows matching the given bookmark filter.
|
||||
/// </summary>
|
||||
Task<IEnumerable<RunWorkflowInstanceResponse>> ResumeAsync(BookmarkFilter filter, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -276,6 +276,7 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module)
|
|||
.AddScoped<IBookmarksPersister, BookmarksPersister>()
|
||||
.AddScoped<IBookmarkResumer, BookmarkResumer>()
|
||||
.AddScoped<IBookmarkQueue, StoreBookmarkQueue>()
|
||||
.AddScoped<IWorkflowResumer, WorkflowResumer>()
|
||||
.AddScoped<ITriggerInvoker, TriggerInvoker>()
|
||||
.AddScoped<IWorkflowCanceler, WorkflowCanceler>()
|
||||
.AddScoped<IWorkflowCancellationService, WorkflowCancellationService>()
|
||||
|
|
|
|||
|
|
@ -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;
|
|||
/// </summary>
|
||||
public class BookmarkFilter
|
||||
{
|
||||
// Cache the properties of BookmarkFilter for performance.
|
||||
private static readonly System.Reflection.PropertyInfo[] CachedProperties = typeof(BookmarkFilter).GetProperties();
|
||||
|
||||
/// <summary>
|
||||
/// Gets or sets the ID of the bookmark.
|
||||
/// </summary>
|
||||
|
|
@ -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<string>();
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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; }
|
||||
|
||||
/// <summary>
|
||||
/// The ID of the activity instance to resume, if any.
|
||||
/// </summary>
|
||||
public string? ActivityInstanceId { get; set; }
|
||||
|
||||
/// Any additional properties to associate with the workflow instance.
|
||||
public IDictionary<string, object>? Properties { get; set; }
|
||||
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ using Elsa.Workflows.Runtime.Options;
|
|||
namespace Elsa.Workflows.Runtime;
|
||||
|
||||
/// <inheritdoc />
|
||||
[Obsolete("Will be removed in a future version.")]
|
||||
public class BookmarkBoundWorkflowService(IWorkflowMatcher workflowMatcher) : IBookmarkBoundWorkflowService
|
||||
{
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
|
|
@ -7,7 +7,7 @@ using Microsoft.Extensions.Logging;
|
|||
|
||||
namespace Elsa.Workflows.Runtime;
|
||||
|
||||
public class BookmarkQueueProcessor(IBookmarkQueueStore store, IBookmarkResumer bookmarkResumer, ILogger<BookmarkQueueProcessor> logger) : IBookmarkQueueProcessor
|
||||
public class BookmarkQueueProcessor(IBookmarkQueueStore store, IWorkflowResumer workflowResumer, ILogger<BookmarkQueueProcessor> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -8,6 +8,7 @@ using Microsoft.Extensions.Logging;
|
|||
namespace Elsa.Workflows.Runtime;
|
||||
|
||||
/// <inheritdoc />
|
||||
[Obsolete("Use WorkflowResumer instead.")]
|
||||
public class BookmarkResumer(IWorkflowRuntime workflowRuntime, IBookmarkStore bookmarkStore, IStimulusHasher stimulusHasher, ILogger<BookmarkResumer> logger) : IBookmarkResumer
|
||||
{
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
|
|
@ -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<StimulusSender> 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<ICollection<RunWorkflowInstanceResponse>> 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<RunWorkflowInstanceResponse>();
|
||||
|
||||
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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
135
src/modules/Elsa.Workflows.Runtime/Services/WorkflowResumer.cs
Normal file
135
src/modules/Elsa.Workflows.Runtime/Services/WorkflowResumer.cs
Normal file
|
|
@ -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;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class WorkflowResumer(
|
||||
IWorkflowRuntime workflowRuntime,
|
||||
IBookmarkStore bookmarkStore,
|
||||
IStimulusHasher stimulusHasher,
|
||||
IDistributedLockProvider distributedLockProvider,
|
||||
IOptions<DistributedLockingOptions> distributedLockingOptions,
|
||||
ILogger<WorkflowResumer> logger) : IWorkflowResumer
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public Task<IEnumerable<RunWorkflowInstanceResponse>> ResumeAsync<TActivity>(object stimulus, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity
|
||||
{
|
||||
return ResumeAsync<TActivity>(stimulus, null, options, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<IEnumerable<RunWorkflowInstanceResponse>> ResumeAsync<TActivity>(object stimulus, string? workflowInstanceId = null, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity
|
||||
{
|
||||
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<TActivity>();
|
||||
var stimulusHash = stimulusHasher.Hash(activityTypeName, stimulus);
|
||||
var bookmarkFilter = new BookmarkFilter
|
||||
{
|
||||
Name = activityTypeName,
|
||||
WorkflowInstanceId = workflowInstanceId,
|
||||
Hash = stimulusHash,
|
||||
};
|
||||
return await ResumeAsync(bookmarkFilter, options, cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<RunWorkflowInstanceResponse?> ResumeAsync(string bookmarkId, IDictionary<string, object> 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();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<RunWorkflowInstanceResponse?> ResumeAsync<TActivity>(string bookmarkId, ResumeBookmarkOptions? options = null, CancellationToken cancellationToken = default) where TActivity : IActivity
|
||||
{
|
||||
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<TActivity>();
|
||||
var bookmarkFilter = new BookmarkFilter
|
||||
{
|
||||
Name = activityTypeName,
|
||||
BookmarkId = bookmarkId
|
||||
};
|
||||
var response = await ResumeAsync(bookmarkFilter, options, cancellationToken);
|
||||
return response.FirstOrDefault();
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<RunWorkflowInstanceResponse>> 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);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<IEnumerable<RunWorkflowInstanceResponse>> 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<RunWorkflowInstanceResponse>();
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue