Ensure BookmarksDeleting and BookmarksDeleted events are published

Fixes #4389
This commit is contained in:
Sipke Schoorstra 2023-09-10 16:39:42 +02:00
parent 4f73419573
commit cbf609760e
4 changed files with 19 additions and 10 deletions

View file

@ -53,6 +53,7 @@ public class DapperBookmarkStore : IBookmarkStore
{
query
.Is(nameof(StoredBookmarkRecord.Hash), filter.Hash)
.In(nameof(StoredBookmarkRecord.Hash), filter.Hashes)
.Is(nameof(StoredBookmarkRecord.WorkflowInstanceId), filter.WorkflowInstanceId)
.In(nameof(StoredBookmarkRecord.WorkflowInstanceId), filter.WorkflowInstanceIds)
.Is(nameof(StoredBookmarkRecord.CorrelationId), filter.CorrelationId)

View file

@ -35,6 +35,7 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime
private readonly IWorkflowInstanceStore _workflowInstanceStore;
private readonly IIdentityGenerator _identityGenerator;
private readonly IBookmarkHasher _hasher;
private readonly IBookmarkManager _bookmarkManager;
private readonly IWorkflowDefinitionService _workflowDefinitionService;
private readonly IWorkflowInstanceFactory _workflowInstanceFactory;
private readonly WorkflowExecutionResultMapper _workflowExecutionResultMapper;
@ -50,6 +51,7 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime
IWorkflowInstanceStore workflowInstanceStore,
IIdentityGenerator identityGenerator,
IBookmarkHasher hasher,
IBookmarkManager bookmarkManager,
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowInstanceFactory workflowInstanceFactory,
WorkflowExecutionResultMapper workflowExecutionResultMapper)
@ -61,6 +63,7 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime
_workflowInstanceStore = workflowInstanceStore;
_identityGenerator = identityGenerator;
_hasher = hasher;
_bookmarkManager = bookmarkManager;
_workflowDefinitionService = workflowDefinitionService;
_workflowInstanceFactory = workflowInstanceFactory;
_workflowExecutionResultMapper = workflowExecutionResultMapper;
@ -300,11 +303,9 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime
private async Task RemoveBookmarksAsync(string workflowInstanceId, IEnumerable<Bookmark> bookmarks, CancellationToken cancellationToken = default)
{
foreach (var bookmark in bookmarks)
{
var filter = new BookmarkFilter { Hash = bookmark.Hash, WorkflowInstanceId = workflowInstanceId };
await _bookmarkStore.DeleteAsync(filter, cancellationToken);
}
var matchingHashes = bookmarks.Select(x => x.Hash).ToList();
var filter = new BookmarkFilter { Hashes = matchingHashes, WorkflowInstanceId = workflowInstanceId };
await _bookmarkManager.DeleteManyAsync(filter, cancellationToken);
}
private async Task<IEnumerable<WorkflowMatch>> FindStartableWorkflowsAsync(WorkflowsFilter workflowsFilter, CancellationToken cancellationToken)

View file

@ -32,6 +32,11 @@ public class BookmarkFilter
/// </summary>
public string? Hash { get; set; }
/// <summary>
/// Gets or sets the hashes of the bookmarks to find.
/// </summary>
public ICollection<string>? Hashes { get; set; }
/// <summary>
/// Gets or sets the correlation ID of the bookmark to find.
/// </summary>
@ -62,6 +67,7 @@ public class BookmarkFilter
if (filter.BookmarkIds != null) query = query.Where(x => filter.BookmarkIds.Contains(x.BookmarkId));
if (filter.CorrelationId != null) query = query.Where(x => x.CorrelationId == filter.CorrelationId);
if (filter.Hash != null) query = query.Where(x => x.Hash == filter.Hash);
if (filter.Hashes != null) query = query.Where(x => filter.Hashes.Contains(x.Hash));
if (filter.WorkflowInstanceId != null) query = query.Where(x => x.WorkflowInstanceId == filter.WorkflowInstanceId);
if (filter.WorkflowInstanceIds != null) query = query.Where(x => filter.WorkflowInstanceIds.Contains(x.WorkflowInstanceId));
if (filter.ActivityTypeName != null) query = query.Where(x => x.ActivityTypeName == filter.ActivityTypeName);

View file

@ -27,6 +27,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
private readonly ITriggerStore _triggerStore;
private readonly IBookmarkStore _bookmarkStore;
private readonly IBookmarkHasher _hasher;
private readonly IBookmarkManager _bookmarkManager;
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly IWorkflowInstanceFactory _workflowInstanceFactory;
private readonly WorkflowStateMapper _workflowStateMapper;
@ -42,6 +43,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
ITriggerStore triggerStore,
IBookmarkStore bookmarkStore,
IBookmarkHasher hasher,
IBookmarkManager bookmarkManager,
IDistributedLockProvider distributedLockProvider,
IWorkflowInstanceFactory workflowInstanceFactory,
WorkflowStateMapper workflowStateMapper,
@ -53,6 +55,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
_triggerStore = triggerStore;
_bookmarkStore = bookmarkStore;
_hasher = hasher;
_bookmarkManager = bookmarkManager;
_distributedLockProvider = distributedLockProvider;
_workflowInstanceFactory = workflowInstanceFactory;
_workflowStateMapper = workflowStateMapper;
@ -351,11 +354,9 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
private async Task RemoveBookmarksAsync(string workflowInstanceId, IEnumerable<Bookmark> bookmarks, CancellationToken cancellationToken)
{
foreach (var bookmark in bookmarks)
{
var filter = new BookmarkFilter { Hash = bookmark.Hash, WorkflowInstanceId = workflowInstanceId };
await _bookmarkStore.DeleteAsync(filter, cancellationToken);
}
var matchingHashes = bookmarks.Select(x => x.Hash).ToList();
var filter = new BookmarkFilter { Hashes = matchingHashes, WorkflowInstanceId = workflowInstanceId };
await _bookmarkManager.DeleteManyAsync(filter, cancellationToken);
}
private async Task<IEnumerable<WorkflowMatch>> FindStartableWorkflowsAsync(