| | | 1 | | using Elsa.Workflows; |
| | | 2 | | using Elsa.Workflows.Management; |
| | | 3 | | using Elsa.Workflows.Management.Filters; |
| | | 4 | | using Elsa.Workflows.Runtime; |
| | | 5 | | using Elsa.Workflows.Runtime.Entities; |
| | | 6 | | using Elsa.Workflows.Runtime.Filters; |
| | | 7 | | |
| | | 8 | | namespace Elsa.Scheduling.Services; |
| | | 9 | | |
| | | 10 | | /// <summary> |
| | | 11 | | /// Classifies Delay/Timer/Cron/StartAt stored bookmarks against the workflow-instance store so startup |
| | | 12 | | /// schedule rebuild can skip and purge bookmarks whose instance is missing or finished. |
| | | 13 | | /// </summary> |
| | | 14 | | public class SchedulingBookmarkReconciler(IWorkflowInstanceStore workflowInstanceStore, IBookmarkManager? bookmarkManage |
| | | 15 | | { |
| | | 16 | | /// <summary> |
| | | 17 | | /// Splits bookmarks into those that may be scheduled and those whose instance is missing or terminal. |
| | | 18 | | /// </summary> |
| | | 19 | | public async Task<SchedulingBookmarkClassification> ClassifyAsync(IEnumerable<StoredBookmark> bookmarks, Cancellatio |
| | | 20 | | { |
| | | 21 | | var bookmarkList = bookmarks as IReadOnlyList<StoredBookmark> ?? bookmarks.ToList(); |
| | | 22 | | |
| | | 23 | | if (bookmarkList.Count == 0) |
| | | 24 | | return new SchedulingBookmarkClassification([], []); |
| | | 25 | | |
| | | 26 | | var instanceIds = bookmarkList |
| | | 27 | | .Select(x => x.WorkflowInstanceId) |
| | | 28 | | .Where(x => !string.IsNullOrWhiteSpace(x)) |
| | | 29 | | .Distinct(StringComparer.Ordinal) |
| | | 30 | | .ToList(); |
| | | 31 | | |
| | | 32 | | var liveInstanceIds = instanceIds.Count == 0 |
| | | 33 | | ? new HashSet<string>(StringComparer.Ordinal) |
| | | 34 | | : (await workflowInstanceStore.FindManyIdsAsync(new WorkflowInstanceFilter |
| | | 35 | | { |
| | | 36 | | Ids = instanceIds, |
| | | 37 | | WorkflowStatus = WorkflowStatus.Running |
| | | 38 | | }, cancellationToken)).ToHashSet(StringComparer.Ordinal); |
| | | 39 | | |
| | | 40 | | var schedulable = new List<StoredBookmark>(bookmarkList.Count); |
| | | 41 | | var orphans = new List<StoredBookmark>(); |
| | | 42 | | |
| | | 43 | | foreach (var bookmark in bookmarkList) |
| | | 44 | | { |
| | | 45 | | if (!string.IsNullOrWhiteSpace(bookmark.WorkflowInstanceId) && liveInstanceIds.Contains(bookmark.WorkflowIns |
| | | 46 | | schedulable.Add(bookmark); |
| | | 47 | | else |
| | | 48 | | orphans.Add(bookmark); |
| | | 49 | | } |
| | | 50 | | |
| | | 51 | | return new SchedulingBookmarkClassification(schedulable, orphans); |
| | | 52 | | } |
| | | 53 | | |
| | | 54 | | /// <summary> |
| | | 55 | | /// Revalidates candidates against the workflow-instance store, then deletes bookmarks that are still orphans. |
| | | 56 | | /// </summary> |
| | | 57 | | public async Task PurgeAsync(IEnumerable<StoredBookmark> candidateBookmarks, CancellationToken cancellationToken = d |
| | | 58 | | { |
| | | 59 | | if (bookmarkManager == null) |
| | | 60 | | return; |
| | | 61 | | |
| | | 62 | | var remainingOrphans = (await ClassifyAsync(candidateBookmarks, cancellationToken)).Orphans; |
| | | 63 | | var bookmarkIds = remainingOrphans |
| | | 64 | | .Select(x => x.Id) |
| | | 65 | | .Where(x => !string.IsNullOrWhiteSpace(x)) |
| | | 66 | | .Distinct(StringComparer.Ordinal) |
| | | 67 | | .ToList(); |
| | | 68 | | |
| | | 69 | | if (bookmarkIds.Count == 0) |
| | | 70 | | return; |
| | | 71 | | |
| | | 72 | | await bookmarkManager.DeleteManyAsync(new BookmarkFilter |
| | | 73 | | { |
| | | 74 | | BookmarkIds = bookmarkIds |
| | | 75 | | }, cancellationToken); |
| | | 76 | | } |
| | | 77 | | } |
| | | 78 | | |
| | | 79 | | /// <summary> |
| | | 80 | | /// The schedulable vs orphan split for a batch of scheduling bookmarks. |
| | | 81 | | /// </summary> |
| | 73 | 82 | | public record SchedulingBookmarkClassification(IReadOnlyList<StoredBookmark> Schedulable, IReadOnlyList<StoredBookmark> |