| | | 1 | | using Elsa.Common; |
| | | 2 | | using Elsa.Common.Multitenancy; |
| | | 3 | | using Elsa.Common.Models; |
| | | 4 | | using Elsa.Scheduling.Options; |
| | | 5 | | using Elsa.Scheduling.Services; |
| | | 6 | | using Elsa.Workflows.Management; |
| | | 7 | | using Elsa.Workflows.Runtime; |
| | | 8 | | using Elsa.Workflows.Runtime.Filters; |
| | | 9 | | using Elsa.Workflows.Runtime.Tasks; |
| | | 10 | | using Microsoft.Extensions.DependencyInjection; |
| | | 11 | | using Microsoft.Extensions.Options; |
| | | 12 | | |
| | | 13 | | namespace Elsa.Scheduling.StartupTasks; |
| | | 14 | | |
| | | 15 | | /// <summary> |
| | | 16 | | /// Enqueues schedule creation when using the default scheduler, which doesn't have its own persistence layer like Quart |
| | | 17 | | /// Scheduling bookmarks whose workflow instance is missing or finished are skipped and purged so startup does not rehyd |
| | | 18 | | /// </summary> |
| | | 19 | | [TaskDependency(typeof(PopulateRegistriesStartupTask))] |
| | 100 | 20 | | public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptions<SchedulingOptions> options) : IStartu |
| | | 21 | | { |
| | | 22 | | public async Task ExecuteAsync(CancellationToken cancellationToken) |
| | | 23 | | { |
| | 100 | 24 | | var workQueue = serviceProvider.GetService<ITenantBackgroundWorkQueue>(); |
| | | 25 | | |
| | 100 | 26 | | if (workQueue != null) |
| | 92 | 27 | | await workQueue.EnqueueAsync(CreateSchedulesAsync, cancellationToken); |
| | | 28 | | else |
| | 8 | 29 | | await CreateSchedulesAsync(serviceProvider, cancellationToken); |
| | 100 | 30 | | } |
| | | 31 | | |
| | | 32 | | private async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) |
| | | 33 | | { |
| | 100 | 34 | | var triggerStore = serviceProvider.GetRequiredService<ITriggerStore>(); |
| | 100 | 35 | | var bookmarkStore = serviceProvider.GetRequiredService<IBookmarkStore>(); |
| | 100 | 36 | | var triggerScheduler = serviceProvider.GetRequiredService<ITriggerScheduler>(); |
| | 100 | 37 | | var bookmarkScheduler = serviceProvider.GetRequiredService<IBookmarkScheduler>(); |
| | 100 | 38 | | var workflowInstanceStore = serviceProvider.GetService<IWorkflowInstanceStore>(); |
| | 100 | 39 | | var bookmarkReconciler = workflowInstanceStore == null |
| | 100 | 40 | | ? null |
| | 100 | 41 | | : new SchedulingBookmarkReconciler(workflowInstanceStore, serviceProvider.GetService<IBookmarkManager>()); |
| | 100 | 42 | | var pageSize = Math.Max(1, options.Value.StartupSchedulePageSize); |
| | 100 | 43 | | var stimulusNames = new[] |
| | 100 | 44 | | { |
| | 100 | 45 | | SchedulingStimulusNames.Cron, SchedulingStimulusNames.Timer, SchedulingStimulusNames.StartAt, SchedulingStim |
| | 100 | 46 | | }; |
| | 100 | 47 | | var triggerFilter = new TriggerFilter |
| | 100 | 48 | | { |
| | 100 | 49 | | Names = stimulusNames |
| | 100 | 50 | | }; |
| | 100 | 51 | | var bookmarkFilter = new BookmarkFilter |
| | 100 | 52 | | { |
| | 100 | 53 | | Names = stimulusNames |
| | 100 | 54 | | }; |
| | | 55 | | |
| | 100 | 56 | | await ScheduleTriggersAsync(triggerStore, triggerScheduler, triggerFilter, pageSize, cancellationToken); |
| | 100 | 57 | | await ScheduleBookmarksAsync(bookmarkStore, bookmarkScheduler, bookmarkReconciler, bookmarkFilter, pageSize, can |
| | 100 | 58 | | } |
| | | 59 | | |
| | | 60 | | private static async Task ScheduleTriggersAsync(ITriggerStore triggerStore, ITriggerScheduler triggerScheduler, Trig |
| | | 61 | | { |
| | 100 | 62 | | var pageArgs = PageArgs.FromRange(0, pageSize); |
| | | 63 | | |
| | 1 | 64 | | while (true) |
| | | 65 | | { |
| | 101 | 66 | | var page = await triggerStore.FindManyAsync(triggerFilter, pageArgs, cancellationToken); |
| | | 67 | | |
| | 101 | 68 | | if (page.Items.Count == 0) |
| | | 69 | | break; |
| | | 70 | | |
| | 10 | 71 | | await triggerScheduler.ScheduleAsync(page.Items, cancellationToken); |
| | | 72 | | |
| | 10 | 73 | | var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count; |
| | 10 | 74 | | if (nextOffset >= page.TotalCount) |
| | | 75 | | break; |
| | | 76 | | |
| | 1 | 77 | | pageArgs = pageArgs.Next(); |
| | 1 | 78 | | } |
| | 100 | 79 | | } |
| | | 80 | | |
| | | 81 | | private static async Task ScheduleBookmarksAsync( |
| | | 82 | | IBookmarkStore bookmarkStore, |
| | | 83 | | IBookmarkScheduler bookmarkScheduler, |
| | | 84 | | SchedulingBookmarkReconciler? bookmarkReconciler, |
| | | 85 | | BookmarkFilter bookmarkFilter, |
| | | 86 | | int pageSize, |
| | | 87 | | CancellationToken cancellationToken) |
| | | 88 | | { |
| | 100 | 89 | | var pageArgs = PageArgs.FromRange(0, pageSize); |
| | 100 | 90 | | var orphanBookmarkIds = new HashSet<string>(StringComparer.Ordinal); |
| | | 91 | | |
| | 2 | 92 | | while (true) |
| | | 93 | | { |
| | 102 | 94 | | var page = await bookmarkStore.FindManyAsync(bookmarkFilter, pageArgs, cancellationToken); |
| | | 95 | | |
| | 102 | 96 | | if (page.Items.Count == 0) |
| | | 97 | | break; |
| | | 98 | | |
| | 11 | 99 | | if (bookmarkReconciler == null) |
| | | 100 | | { |
| | 1 | 101 | | await bookmarkScheduler.ScheduleAsync(page.Items, cancellationToken); |
| | | 102 | | } |
| | | 103 | | else |
| | | 104 | | { |
| | 10 | 105 | | var classification = await bookmarkReconciler.ClassifyAsync(page.Items, cancellationToken); |
| | 30 | 106 | | orphanBookmarkIds.UnionWith(classification.Orphans.Select(x => x.Id).Where(id => !string.IsNullOrWhiteSp |
| | | 107 | | |
| | 10 | 108 | | if (classification.Schedulable.Count > 0) |
| | 7 | 109 | | await bookmarkScheduler.ScheduleAsync(classification.Schedulable, cancellationToken); |
| | | 110 | | } |
| | | 111 | | |
| | 11 | 112 | | var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count; |
| | 11 | 113 | | if (nextOffset >= page.TotalCount) |
| | | 114 | | break; |
| | | 115 | | |
| | 2 | 116 | | pageArgs = pageArgs.Next(); |
| | 2 | 117 | | } |
| | | 118 | | |
| | 100 | 119 | | if (bookmarkReconciler == null || orphanBookmarkIds.Count == 0) |
| | 96 | 120 | | return; |
| | | 121 | | |
| | 18 | 122 | | foreach (var orphanIdBatch in orphanBookmarkIds.Chunk(pageSize)) |
| | | 123 | | { |
| | 5 | 124 | | var candidates = await bookmarkStore.FindManyAsync(new BookmarkFilter |
| | 5 | 125 | | { |
| | 5 | 126 | | BookmarkIds = orphanIdBatch.ToList() |
| | 5 | 127 | | }, cancellationToken); |
| | 5 | 128 | | var classification = await bookmarkReconciler.ClassifyAsync(candidates, cancellationToken); |
| | | 129 | | |
| | 5 | 130 | | if (classification.Schedulable.Count > 0) |
| | 1 | 131 | | await bookmarkScheduler.ScheduleAsync(classification.Schedulable, cancellationToken); |
| | | 132 | | |
| | 5 | 133 | | await bookmarkReconciler.PurgeAsync(classification.Orphans, cancellationToken); |
| | 5 | 134 | | } |
| | 100 | 135 | | } |
| | | 136 | | } |