using Elsa.Common; using Elsa.Common.Multitenancy; using Elsa.Common.Models; using Elsa.Scheduling.Options; using Elsa.Scheduling.StartupTasks; using Elsa.Workflows.Management; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Runtime; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.Tasks; using Microsoft.Extensions.DependencyInjection; using NSubstitute; using OptionsFactory = Microsoft.Extensions.Options.Options; namespace Elsa.Scheduling.UnitTests.StartupTasks; public class CreateSchedulesStartupTaskTests { private readonly StoredTrigger[] _triggers = [ new() { Id = "trigger-1", WorkflowDefinitionId = "definition", WorkflowDefinitionVersionId = "version", ActivityId = "activity" } ]; private readonly StoredBookmark[] _bookmarks = [new() { Id = "bookmark-1", Hash = "hash", WorkflowInstanceId = "instance" }]; private readonly ITriggerStore _triggerStore = Substitute.For(); private readonly IBookmarkStore _bookmarkStore = Substitute.For(); private readonly ITriggerScheduler _triggerScheduler = Substitute.For(); private readonly IBookmarkScheduler _bookmarkScheduler = Substitute.For(); private readonly IWorkflowInstanceStore _workflowInstanceStore = Substitute.For(); private readonly IBookmarkManager _bookmarkManager = Substitute.For(); private readonly SchedulingOptions _options = new() { StartupSchedulePageSize = 1000 }; public CreateSchedulesStartupTaskTests() { _triggerStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(new Page(_triggers, _triggers.Length)); _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(new Page(_bookmarks, _bookmarks.Length)); _workflowInstanceStore.FindManyIdsAsync(Arg.Any(), Arg.Any()) .Returns(call => RunningIds(call.Arg())); } [Fact] public void Task_DependsOnPopulateRegistriesStartupTask() { var dependency = Assert.Single(typeof(CreateSchedulesStartupTask).GetCustomAttributes(typeof(TaskDependencyAttribute), false).Cast()); Assert.Equal(typeof(PopulateRegistriesStartupTask), dependency.DependencyTaskType); } [Fact] public async Task ExecuteAsync_WithoutTenantBackgroundQueue_SchedulesImmediately() { var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(_triggers)), Arg.Any()); await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(_bookmarks)), Arg.Any()); } [Fact] public async Task ExecuteAsync_WithTenantBackgroundQueue_EnqueuesScheduleCreation() { TenantBackgroundWorkItem? workItem = null; var workQueue = Substitute.For(); workQueue.EnqueueAsync(Arg.Do(x => workItem = x), Arg.Any()).Returns(ValueTask.CompletedTask); var serviceProvider = CreateServiceProvider(services => services.AddSingleton(workQueue)); var task = new CreateSchedulesStartupTask(serviceProvider, OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await workQueue.Received(1).EnqueueAsync(Arg.Any(), Arg.Any()); await _triggerScheduler.DidNotReceive().ScheduleAsync(Arg.Any>(), Arg.Any()); Assert.NotNull(workItem); await workItem(serviceProvider, CancellationToken.None); await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(_triggers)), Arg.Any()); await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(_bookmarks)), Arg.Any()); } [Fact] public async Task ExecuteAsync_SchedulesInConfiguredPages() { var firstTriggerPage = new[] { _triggers[0] }; var secondTriggerPage = new[] { new StoredTrigger { Id = "trigger-2", WorkflowDefinitionId = "definition", WorkflowDefinitionVersionId = "version", ActivityId = "activity" } }; var firstBookmarkPage = new[] { _bookmarks[0] }; var secondBookmarkPage = new[] { new StoredBookmark { Id = "bookmark-2", Hash = "hash", WorkflowInstanceId = "instance" } }; _options.StartupSchedulePageSize = 1; _triggerStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 0 && x.Limit == 1), Arg.Any()) .Returns(new Page(firstTriggerPage, 2)); _triggerStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 1 && x.Limit == 1), Arg.Any()) .Returns(new Page(secondTriggerPage, 2)); _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 0 && x.Limit == 1), Arg.Any()) .Returns(new Page(firstBookmarkPage, 2)); _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 1 && x.Limit == 1), Arg.Any()) .Returns(new Page(secondBookmarkPage, 2)); var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(firstTriggerPage)), Arg.Any()); await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(secondTriggerPage)), Arg.Any()); await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(firstBookmarkPage)), Arg.Any()); await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(secondBookmarkPage)), Arg.Any()); } [Fact] public async Task ExecuteAsync_SkipsAndPurgesOrphanSchedulingBookmarks() { var missingInstance = Bookmark("missing-bookmark", "missing-instance"); var finishedInstance = Bookmark("finished-bookmark", "finished-instance"); var cancelledInstance = Bookmark("cancelled-bookmark", "cancelled-instance"); var faultedInstance = Bookmark("faulted-bookmark", "faulted-instance"); var emptyInstanceId = Bookmark("empty-instance-bookmark", ""); var suspendedInstance = Bookmark("suspended-bookmark", "suspended-instance"); var bookmarks = new[] { missingInstance, finishedInstance, cancelledInstance, faultedInstance, emptyInstanceId, suspendedInstance }; BookmarkFilter? purgedFilter = null; _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(new Page(bookmarks, bookmarks.Length)); StubBookmarkReload(bookmarks); _workflowInstanceStore.FindManyIdsAsync(Arg.Any(), Arg.Any()) .Returns(["suspended-instance"]); _bookmarkManager.DeleteManyAsync(Arg.Do(x => purgedFilter = x), Arg.Any()) .Returns(5); var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _bookmarkScheduler.Received(1).ScheduleAsync( Arg.Is>(x => x.SequenceEqual(new[] { suspendedInstance })), Arg.Any()); await _bookmarkManager.Received(1).DeleteManyAsync(Arg.Any(), Arg.Any()); Assert.NotNull(purgedFilter); Assert.Equal( new[] { "cancelled-bookmark", "empty-instance-bookmark", "faulted-bookmark", "finished-bookmark", "missing-bookmark" }, purgedFilter.BookmarkIds?.OrderBy(x => x, StringComparer.Ordinal).ToArray()); } [Fact] public async Task ExecuteAsync_SchedulesOrphanWhoseInstanceBecomesRunningBeforePurge() { var revived = Bookmark("revived-bookmark", "revived-instance"); var stillMissing = Bookmark("missing-bookmark", "missing-instance"); var bookmarks = new[] { revived, stillMissing }; BookmarkFilter? purgedFilter = null; _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(new Page(bookmarks, bookmarks.Length)); StubBookmarkReload(bookmarks); _workflowInstanceStore.FindManyIdsAsync(Arg.Any(), Arg.Any()) .Returns(_ => Array.Empty(), _ => new[] { "revived-instance" }); _bookmarkManager.DeleteManyAsync(Arg.Do(x => purgedFilter = x), Arg.Any()) .Returns(1); var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _bookmarkScheduler.Received(1).ScheduleAsync( Arg.Is>(x => x.SequenceEqual(new[] { revived })), Arg.Any()); await _bookmarkManager.Received(1).DeleteManyAsync(Arg.Any(), Arg.Any()); Assert.NotNull(purgedFilter); Assert.Equal(["missing-bookmark"], purgedFilter.BookmarkIds); } [Fact] public async Task ExecuteAsync_ReconcilesOrphansInConfiguredPages() { var firstOrphan = Bookmark("orphan-1", "missing-1"); var secondOrphan = Bookmark("orphan-2", "missing-2"); var bookmarks = new[] { firstOrphan, secondOrphan }; var purgedFilters = new List(); var reloadedIdBatches = new List>(); _options.StartupSchedulePageSize = 1; _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 0 && x.Limit == 1), Arg.Any()) .Returns(new Page([firstOrphan], 2)); _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Is(x => x.Offset == 1 && x.Limit == 1), Arg.Any()) .Returns(new Page([secondOrphan], 2)); _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any()) .Returns(call => { var ids = call.Arg().BookmarkIds ?? []; reloadedIdBatches.Add(ids.ToList()); return bookmarks.Where(x => ids.Contains(x.Id)); }); _workflowInstanceStore.FindManyIdsAsync(Arg.Any(), Arg.Any()) .Returns(Array.Empty()); _bookmarkManager.DeleteManyAsync(Arg.Do(x => purgedFilters.Add(x)), Arg.Any()) .Returns(1); var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _bookmarkScheduler.DidNotReceive().ScheduleAsync(Arg.Any>(), Arg.Any()); Assert.Equal(2, reloadedIdBatches.Count); Assert.All(reloadedIdBatches, batch => Assert.Single(batch)); Assert.Equal(2, purgedFilters.Count); Assert.All(purgedFilters, filter => Assert.Single(filter.BookmarkIds!)); Assert.Equal(["orphan-1", "orphan-2"], purgedFilters.SelectMany(x => x.BookmarkIds!).OrderBy(x => x, StringComparer.Ordinal)); } [Fact] public async Task ExecuteAsync_WithoutWorkflowInstanceStore_SchedulesAllBookmarks() { var missingInstance = Bookmark("missing-bookmark", "missing-instance"); var bookmarks = new[] { missingInstance, _bookmarks[0] }; _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(new Page(bookmarks, bookmarks.Length)); var task = new CreateSchedulesStartupTask(CreateServiceProvider(registerReconcileServices: false), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(bookmarks)), Arg.Any()); await _bookmarkManager.DidNotReceive().DeleteManyAsync(Arg.Any(), Arg.Any()); } [Fact] public async Task ExecuteAsync_WithoutBookmarkManager_SkipsOrphansButDoesNotPurge() { var missingInstance = Bookmark("missing-bookmark", "missing-instance"); var suspendedInstance = Bookmark("suspended-bookmark", "suspended-instance"); var bookmarks = new[] { missingInstance, suspendedInstance }; _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(new Page(bookmarks, bookmarks.Length)); StubBookmarkReload(bookmarks); _workflowInstanceStore.FindManyIdsAsync(Arg.Any(), Arg.Any()) .Returns(["suspended-instance"]); var task = new CreateSchedulesStartupTask( CreateServiceProvider(services => services.AddSingleton(_workflowInstanceStore), registerReconcileServices: false), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _bookmarkScheduler.Received(1).ScheduleAsync( Arg.Is>(x => x.SequenceEqual(new[] { suspendedInstance })), Arg.Any()); await _bookmarkManager.DidNotReceive().DeleteManyAsync(Arg.Any(), Arg.Any()); } [Fact] public async Task ExecuteAsync_DoesNotPurgeWhenEverySchedulingBookmarkIsLive() { var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options)); await task.ExecuteAsync(CancellationToken.None); await _bookmarkManager.DidNotReceive().DeleteManyAsync(Arg.Any(), Arg.Any()); } private ServiceProvider CreateServiceProvider(Action? configureServices = null, bool registerReconcileServices = true) { var services = new ServiceCollection(); services.AddSingleton(_triggerStore); services.AddSingleton(_bookmarkStore); services.AddSingleton(_triggerScheduler); services.AddSingleton(_bookmarkScheduler); if (registerReconcileServices) { services.AddSingleton(_workflowInstanceStore); services.AddSingleton(_bookmarkManager); } configureServices?.Invoke(services); return services.BuildServiceProvider(); } private void StubBookmarkReload(IEnumerable bookmarks) { var bookmarkList = bookmarks.ToList(); _bookmarkStore.FindManyAsync(Arg.Any(), Arg.Any()) .Returns(call => { var ids = call.Arg().BookmarkIds; if (ids == null) return bookmarkList.AsEnumerable(); var idSet = ids.ToHashSet(StringComparer.Ordinal); return bookmarkList.Where(x => idSet.Contains(x.Id)); }); } private static IEnumerable RunningIds(WorkflowInstanceFilter filter) { return (filter.Ids ?? []).Where(id => !string.IsNullOrWhiteSpace(id)); } private static StoredBookmark Bookmark(string id, string workflowInstanceId) => new() { Id = id, Hash = "hash", Name = SchedulingStimulusNames.Delay, WorkflowInstanceId = workflowInstanceId }; }