| | | 1 | | using Elsa.Common; |
| | | 2 | | using Elsa.Common.Multitenancy; |
| | | 3 | | using Elsa.Common.Models; |
| | | 4 | | using Elsa.Scheduling.Options; |
| | | 5 | | using Elsa.Workflows.Runtime; |
| | | 6 | | using Elsa.Workflows.Runtime.Filters; |
| | | 7 | | using Elsa.Workflows.Runtime.Tasks; |
| | | 8 | | using Microsoft.Extensions.DependencyInjection; |
| | | 9 | | using Microsoft.Extensions.Options; |
| | | 10 | | |
| | | 11 | | namespace Elsa.Scheduling.StartupTasks; |
| | | 12 | | |
| | | 13 | | /// <summary> |
| | | 14 | | /// Enqueues schedule creation when using the default scheduler, which doesn't have its own persistence layer like Quart |
| | | 15 | | /// </summary> |
| | | 16 | | [TaskDependency(typeof(PopulateRegistriesStartupTask))] |
| | 94 | 17 | | public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptions<SchedulingOptions> options) : IStartu |
| | | 18 | | { |
| | | 19 | | public async Task ExecuteAsync(CancellationToken cancellationToken) |
| | | 20 | | { |
| | 94 | 21 | | var workQueue = serviceProvider.GetService<ITenantBackgroundWorkQueue>(); |
| | | 22 | | |
| | 94 | 23 | | if (workQueue != null) |
| | 92 | 24 | | await workQueue.EnqueueAsync(CreateSchedulesAsync, cancellationToken); |
| | | 25 | | else |
| | 2 | 26 | | await CreateSchedulesAsync(serviceProvider, cancellationToken); |
| | 94 | 27 | | } |
| | | 28 | | |
| | | 29 | | private async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) |
| | | 30 | | { |
| | 94 | 31 | | var triggerStore = serviceProvider.GetRequiredService<ITriggerStore>(); |
| | 94 | 32 | | var bookmarkStore = serviceProvider.GetRequiredService<IBookmarkStore>(); |
| | 94 | 33 | | var triggerScheduler = serviceProvider.GetRequiredService<ITriggerScheduler>(); |
| | 94 | 34 | | var bookmarkScheduler = serviceProvider.GetRequiredService<IBookmarkScheduler>(); |
| | 94 | 35 | | var pageSize = Math.Max(1, options.Value.StartupSchedulePageSize); |
| | 94 | 36 | | var stimulusNames = new[] |
| | 94 | 37 | | { |
| | 94 | 38 | | SchedulingStimulusNames.Cron, SchedulingStimulusNames.Timer, SchedulingStimulusNames.StartAt, SchedulingStim |
| | 94 | 39 | | }; |
| | 94 | 40 | | var triggerFilter = new TriggerFilter |
| | 94 | 41 | | { |
| | 94 | 42 | | Names = stimulusNames |
| | 94 | 43 | | }; |
| | 94 | 44 | | var bookmarkFilter = new BookmarkFilter |
| | 94 | 45 | | { |
| | 94 | 46 | | Names = stimulusNames |
| | 94 | 47 | | }; |
| | | 48 | | |
| | 94 | 49 | | await ScheduleTriggersAsync(triggerStore, triggerScheduler, triggerFilter, pageSize, cancellationToken); |
| | 94 | 50 | | await ScheduleBookmarksAsync(bookmarkStore, bookmarkScheduler, bookmarkFilter, pageSize, cancellationToken); |
| | 94 | 51 | | } |
| | | 52 | | |
| | | 53 | | private static async Task ScheduleTriggersAsync(ITriggerStore triggerStore, ITriggerScheduler triggerScheduler, Trig |
| | | 54 | | { |
| | 94 | 55 | | var pageArgs = PageArgs.FromRange(0, pageSize); |
| | | 56 | | |
| | 1 | 57 | | while (true) |
| | | 58 | | { |
| | 95 | 59 | | var page = await triggerStore.FindManyAsync(triggerFilter, pageArgs, cancellationToken); |
| | | 60 | | |
| | 95 | 61 | | if (page.Items.Count == 0) |
| | | 62 | | break; |
| | | 63 | | |
| | 4 | 64 | | await triggerScheduler.ScheduleAsync(page.Items, cancellationToken); |
| | | 65 | | |
| | 4 | 66 | | var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count; |
| | 4 | 67 | | if (nextOffset >= page.TotalCount) |
| | | 68 | | break; |
| | | 69 | | |
| | 1 | 70 | | pageArgs = pageArgs.Next(); |
| | 1 | 71 | | } |
| | 94 | 72 | | } |
| | | 73 | | |
| | | 74 | | private static async Task ScheduleBookmarksAsync(IBookmarkStore bookmarkStore, IBookmarkScheduler bookmarkScheduler, |
| | | 75 | | { |
| | 94 | 76 | | var pageArgs = PageArgs.FromRange(0, pageSize); |
| | | 77 | | |
| | 1 | 78 | | while (true) |
| | | 79 | | { |
| | 95 | 80 | | var page = await bookmarkStore.FindManyAsync(bookmarkFilter, pageArgs, cancellationToken); |
| | | 81 | | |
| | 95 | 82 | | if (page.Items.Count == 0) |
| | | 83 | | break; |
| | | 84 | | |
| | 4 | 85 | | await bookmarkScheduler.ScheduleAsync(page.Items, cancellationToken); |
| | | 86 | | |
| | 4 | 87 | | var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count; |
| | 4 | 88 | | if (nextOffset >= page.TotalCount) |
| | | 89 | | break; |
| | | 90 | | |
| | 1 | 91 | | pageArgs = pageArgs.Next(); |
| | 1 | 92 | | } |
| | 94 | 93 | | } |
| | | 94 | | } |