using Elsa.Common;
using Elsa.Common.Multitenancy;
using Elsa.Common.Models;
using Elsa.Scheduling.Options;
using Elsa.Scheduling.Services;
using Elsa.Workflows.Management;
using Elsa.Workflows.Runtime;
using Elsa.Workflows.Runtime.Filters;
using Elsa.Workflows.Runtime.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
namespace Elsa.Scheduling.StartupTasks;
///
/// Enqueues schedule creation when using the default scheduler, which doesn't have its own persistence layer like Quartz or Hangfire.
/// Scheduling bookmarks whose workflow instance is missing or finished are skipped and purged so startup does not rehydrate dead work.
///
[TaskDependency(typeof(PopulateRegistriesStartupTask))]
public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptions options) : IStartupTask
{
public async Task ExecuteAsync(CancellationToken cancellationToken)
{
var workQueue = serviceProvider.GetService();
if (workQueue != null)
await workQueue.EnqueueAsync(CreateSchedulesAsync, cancellationToken);
else
await CreateSchedulesAsync(serviceProvider, cancellationToken);
}
private async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
{
var triggerStore = serviceProvider.GetRequiredService();
var bookmarkStore = serviceProvider.GetRequiredService();
var triggerScheduler = serviceProvider.GetRequiredService();
var bookmarkScheduler = serviceProvider.GetRequiredService();
var workflowInstanceStore = serviceProvider.GetService();
var bookmarkReconciler = workflowInstanceStore == null
? null
: new SchedulingBookmarkReconciler(workflowInstanceStore, serviceProvider.GetService());
var pageSize = Math.Max(1, options.Value.StartupSchedulePageSize);
var stimulusNames = new[]
{
SchedulingStimulusNames.Cron, SchedulingStimulusNames.Timer, SchedulingStimulusNames.StartAt, SchedulingStimulusNames.Delay,
};
var triggerFilter = new TriggerFilter
{
Names = stimulusNames
};
var bookmarkFilter = new BookmarkFilter
{
Names = stimulusNames
};
await ScheduleTriggersAsync(triggerStore, triggerScheduler, triggerFilter, pageSize, cancellationToken);
await ScheduleBookmarksAsync(bookmarkStore, bookmarkScheduler, bookmarkReconciler, bookmarkFilter, pageSize, cancellationToken);
}
private static async Task ScheduleTriggersAsync(ITriggerStore triggerStore, ITriggerScheduler triggerScheduler, TriggerFilter triggerFilter, int pageSize, CancellationToken cancellationToken)
{
var pageArgs = PageArgs.FromRange(0, pageSize);
while (true)
{
var page = await triggerStore.FindManyAsync(triggerFilter, pageArgs, cancellationToken);
if (page.Items.Count == 0)
break;
await triggerScheduler.ScheduleAsync(page.Items, cancellationToken);
var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count;
if (nextOffset >= page.TotalCount)
break;
pageArgs = pageArgs.Next();
}
}
private static async Task ScheduleBookmarksAsync(
IBookmarkStore bookmarkStore,
IBookmarkScheduler bookmarkScheduler,
SchedulingBookmarkReconciler? bookmarkReconciler,
BookmarkFilter bookmarkFilter,
int pageSize,
CancellationToken cancellationToken)
{
var pageArgs = PageArgs.FromRange(0, pageSize);
var orphanBookmarkIds = new HashSet(StringComparer.Ordinal);
while (true)
{
var page = await bookmarkStore.FindManyAsync(bookmarkFilter, pageArgs, cancellationToken);
if (page.Items.Count == 0)
break;
if (bookmarkReconciler == null)
{
await bookmarkScheduler.ScheduleAsync(page.Items, cancellationToken);
}
else
{
var classification = await bookmarkReconciler.ClassifyAsync(page.Items, cancellationToken);
orphanBookmarkIds.UnionWith(classification.Orphans.Select(x => x.Id).Where(id => !string.IsNullOrWhiteSpace(id)));
if (classification.Schedulable.Count > 0)
await bookmarkScheduler.ScheduleAsync(classification.Schedulable, cancellationToken);
}
var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count;
if (nextOffset >= page.TotalCount)
break;
pageArgs = pageArgs.Next();
}
if (bookmarkReconciler == null || orphanBookmarkIds.Count == 0)
return;
foreach (var orphanIdBatch in orphanBookmarkIds.Chunk(pageSize))
{
var candidates = await bookmarkStore.FindManyAsync(new BookmarkFilter
{
BookmarkIds = orphanIdBatch.ToList()
}, cancellationToken);
var classification = await bookmarkReconciler.ClassifyAsync(candidates, cancellationToken);
if (classification.Schedulable.Count > 0)
await bookmarkScheduler.ScheduleAsync(classification.Schedulable, cancellationToken);
await bookmarkReconciler.PurgeAsync(classification.Orphans, cancellationToken);
}
}
}