< Summary

Information
Class: Elsa.Scheduling.StartupTasks.CreateSchedulesStartupTask
Assembly: Elsa.Scheduling
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs
Line coverage
100%
Covered lines: 68
Uncovered lines: 0
Coverable lines: 68
Total lines: 136
Line coverage: 100%
Branch coverage
100%
Covered branches: 24
Total branches: 24
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
ExecuteAsync()100%22100%
CreateSchedulesAsync()100%22100%
ScheduleTriggersAsync()100%44100%
ScheduleBookmarksAsync()100%1616100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs

#LineLine coverage
 1using Elsa.Common;
 2using Elsa.Common.Multitenancy;
 3using Elsa.Common.Models;
 4using Elsa.Scheduling.Options;
 5using Elsa.Scheduling.Services;
 6using Elsa.Workflows.Management;
 7using Elsa.Workflows.Runtime;
 8using Elsa.Workflows.Runtime.Filters;
 9using Elsa.Workflows.Runtime.Tasks;
 10using Microsoft.Extensions.DependencyInjection;
 11using Microsoft.Extensions.Options;
 12
 13namespace 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))]
 10020public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptions<SchedulingOptions> options) : IStartu
 21{
 22    public async Task ExecuteAsync(CancellationToken cancellationToken)
 23    {
 10024        var workQueue = serviceProvider.GetService<ITenantBackgroundWorkQueue>();
 25
 10026        if (workQueue != null)
 9227            await workQueue.EnqueueAsync(CreateSchedulesAsync, cancellationToken);
 28        else
 829            await CreateSchedulesAsync(serviceProvider, cancellationToken);
 10030    }
 31
 32    private async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
 33    {
 10034        var triggerStore = serviceProvider.GetRequiredService<ITriggerStore>();
 10035        var bookmarkStore = serviceProvider.GetRequiredService<IBookmarkStore>();
 10036        var triggerScheduler = serviceProvider.GetRequiredService<ITriggerScheduler>();
 10037        var bookmarkScheduler = serviceProvider.GetRequiredService<IBookmarkScheduler>();
 10038        var workflowInstanceStore = serviceProvider.GetService<IWorkflowInstanceStore>();
 10039        var bookmarkReconciler = workflowInstanceStore == null
 10040            ? null
 10041            : new SchedulingBookmarkReconciler(workflowInstanceStore, serviceProvider.GetService<IBookmarkManager>());
 10042        var pageSize = Math.Max(1, options.Value.StartupSchedulePageSize);
 10043        var stimulusNames = new[]
 10044        {
 10045            SchedulingStimulusNames.Cron, SchedulingStimulusNames.Timer, SchedulingStimulusNames.StartAt, SchedulingStim
 10046        };
 10047        var triggerFilter = new TriggerFilter
 10048        {
 10049            Names = stimulusNames
 10050        };
 10051        var bookmarkFilter = new BookmarkFilter
 10052        {
 10053            Names = stimulusNames
 10054        };
 55
 10056        await ScheduleTriggersAsync(triggerStore, triggerScheduler, triggerFilter, pageSize, cancellationToken);
 10057        await ScheduleBookmarksAsync(bookmarkStore, bookmarkScheduler, bookmarkReconciler, bookmarkFilter, pageSize, can
 10058    }
 59
 60    private static async Task ScheduleTriggersAsync(ITriggerStore triggerStore, ITriggerScheduler triggerScheduler, Trig
 61    {
 10062        var pageArgs = PageArgs.FromRange(0, pageSize);
 63
 164        while (true)
 65        {
 10166            var page = await triggerStore.FindManyAsync(triggerFilter, pageArgs, cancellationToken);
 67
 10168            if (page.Items.Count == 0)
 69                break;
 70
 1071            await triggerScheduler.ScheduleAsync(page.Items, cancellationToken);
 72
 1073            var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count;
 1074            if (nextOffset >= page.TotalCount)
 75                break;
 76
 177            pageArgs = pageArgs.Next();
 178        }
 10079    }
 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    {
 10089        var pageArgs = PageArgs.FromRange(0, pageSize);
 10090        var orphanBookmarkIds = new HashSet<string>(StringComparer.Ordinal);
 91
 292        while (true)
 93        {
 10294            var page = await bookmarkStore.FindManyAsync(bookmarkFilter, pageArgs, cancellationToken);
 95
 10296            if (page.Items.Count == 0)
 97                break;
 98
 1199            if (bookmarkReconciler == null)
 100            {
 1101                await bookmarkScheduler.ScheduleAsync(page.Items, cancellationToken);
 102            }
 103            else
 104            {
 10105                var classification = await bookmarkReconciler.ClassifyAsync(page.Items, cancellationToken);
 30106                orphanBookmarkIds.UnionWith(classification.Orphans.Select(x => x.Id).Where(id => !string.IsNullOrWhiteSp
 107
 10108                if (classification.Schedulable.Count > 0)
 7109                    await bookmarkScheduler.ScheduleAsync(classification.Schedulable, cancellationToken);
 110            }
 111
 11112            var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count;
 11113            if (nextOffset >= page.TotalCount)
 114                break;
 115
 2116            pageArgs = pageArgs.Next();
 2117        }
 118
 100119        if (bookmarkReconciler == null || orphanBookmarkIds.Count == 0)
 96120            return;
 121
 18122        foreach (var orphanIdBatch in orphanBookmarkIds.Chunk(pageSize))
 123        {
 5124            var candidates = await bookmarkStore.FindManyAsync(new BookmarkFilter
 5125            {
 5126                BookmarkIds = orphanIdBatch.ToList()
 5127            }, cancellationToken);
 5128            var classification = await bookmarkReconciler.ClassifyAsync(candidates, cancellationToken);
 129
 5130            if (classification.Schedulable.Count > 0)
 1131                await bookmarkScheduler.ScheduleAsync(classification.Schedulable, cancellationToken);
 132
 5133            await bookmarkReconciler.PurgeAsync(classification.Orphans, cancellationToken);
 5134        }
 100135    }
 136}