diff --git a/src/modules/Elsa.Scheduling/Services/SchedulingBookmarkReconciler.cs b/src/modules/Elsa.Scheduling/Services/SchedulingBookmarkReconciler.cs
new file mode 100644
index 000000000..7fb5259ab
--- /dev/null
+++ b/src/modules/Elsa.Scheduling/Services/SchedulingBookmarkReconciler.cs
@@ -0,0 +1,82 @@
+using Elsa.Workflows;
+using Elsa.Workflows.Management;
+using Elsa.Workflows.Management.Filters;
+using Elsa.Workflows.Runtime;
+using Elsa.Workflows.Runtime.Entities;
+using Elsa.Workflows.Runtime.Filters;
+
+namespace Elsa.Scheduling.Services;
+
+///
+/// Classifies Delay/Timer/Cron/StartAt stored bookmarks against the workflow-instance store so startup
+/// schedule rebuild can skip and purge bookmarks whose instance is missing or finished.
+///
+public class SchedulingBookmarkReconciler(IWorkflowInstanceStore workflowInstanceStore, IBookmarkManager? bookmarkManager)
+{
+ ///
+ /// Splits bookmarks into those that may be scheduled and those whose instance is missing or terminal.
+ ///
+ public async Task ClassifyAsync(IEnumerable bookmarks, CancellationToken cancellationToken = default)
+ {
+ var bookmarkList = bookmarks as IReadOnlyList ?? bookmarks.ToList();
+
+ if (bookmarkList.Count == 0)
+ return new SchedulingBookmarkClassification([], []);
+
+ var instanceIds = bookmarkList
+ .Select(x => x.WorkflowInstanceId)
+ .Where(x => !string.IsNullOrWhiteSpace(x))
+ .Distinct(StringComparer.Ordinal)
+ .ToList();
+
+ var liveInstanceIds = instanceIds.Count == 0
+ ? new HashSet(StringComparer.Ordinal)
+ : (await workflowInstanceStore.FindManyIdsAsync(new WorkflowInstanceFilter
+ {
+ Ids = instanceIds,
+ WorkflowStatus = WorkflowStatus.Running
+ }, cancellationToken)).ToHashSet(StringComparer.Ordinal);
+
+ var schedulable = new List(bookmarkList.Count);
+ var orphans = new List();
+
+ foreach (var bookmark in bookmarkList)
+ {
+ if (!string.IsNullOrWhiteSpace(bookmark.WorkflowInstanceId) && liveInstanceIds.Contains(bookmark.WorkflowInstanceId))
+ schedulable.Add(bookmark);
+ else
+ orphans.Add(bookmark);
+ }
+
+ return new SchedulingBookmarkClassification(schedulable, orphans);
+ }
+
+ ///
+ /// Revalidates candidates against the workflow-instance store, then deletes bookmarks that are still orphans.
+ ///
+ public async Task PurgeAsync(IEnumerable candidateBookmarks, CancellationToken cancellationToken = default)
+ {
+ if (bookmarkManager == null)
+ return;
+
+ var remainingOrphans = (await ClassifyAsync(candidateBookmarks, cancellationToken)).Orphans;
+ var bookmarkIds = remainingOrphans
+ .Select(x => x.Id)
+ .Where(x => !string.IsNullOrWhiteSpace(x))
+ .Distinct(StringComparer.Ordinal)
+ .ToList();
+
+ if (bookmarkIds.Count == 0)
+ return;
+
+ await bookmarkManager.DeleteManyAsync(new BookmarkFilter
+ {
+ BookmarkIds = bookmarkIds
+ }, cancellationToken);
+ }
+}
+
+///
+/// The schedulable vs orphan split for a batch of scheduling bookmarks.
+///
+public record SchedulingBookmarkClassification(IReadOnlyList Schedulable, IReadOnlyList Orphans);
diff --git a/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs b/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs
index 7b4c1caeb..3ba5ca5a1 100644
--- a/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs
+++ b/src/modules/Elsa.Scheduling/StartupTasks/CreateSchedulesStartupTask.cs
@@ -2,6 +2,8 @@ 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;
@@ -12,6 +14,7 @@ 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
@@ -32,6 +35,10 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptio
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[]
{
@@ -47,7 +54,7 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptio
};
await ScheduleTriggersAsync(triggerStore, triggerScheduler, triggerFilter, pageSize, cancellationToken);
- await ScheduleBookmarksAsync(bookmarkStore, bookmarkScheduler, bookmarkFilter, 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)
@@ -71,9 +78,16 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptio
}
}
- private static async Task ScheduleBookmarksAsync(IBookmarkStore bookmarkStore, IBookmarkScheduler bookmarkScheduler, BookmarkFilter bookmarkFilter, int pageSize, CancellationToken cancellationToken)
+ 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)
{
@@ -82,7 +96,18 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptio
if (page.Items.Count == 0)
break;
- await bookmarkScheduler.ScheduleAsync(page.Items, cancellationToken);
+ 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)
@@ -90,5 +115,22 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptio
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);
+ }
}
}
diff --git a/test/unit/Elsa.Scheduling.UnitTests/Services/SchedulingBookmarkReconcilerTests.cs b/test/unit/Elsa.Scheduling.UnitTests/Services/SchedulingBookmarkReconcilerTests.cs
new file mode 100644
index 000000000..90a7912ad
--- /dev/null
+++ b/test/unit/Elsa.Scheduling.UnitTests/Services/SchedulingBookmarkReconcilerTests.cs
@@ -0,0 +1,134 @@
+using Elsa.Common.Services;
+using Elsa.Mediator.Contracts;
+using Elsa.Scheduling.Services;
+using Elsa.Workflows;
+using Elsa.Workflows.Management.Entities;
+using Elsa.Workflows.Management.Stores;
+using Elsa.Workflows.Runtime;
+using Elsa.Workflows.Runtime.Entities;
+using Elsa.Workflows.Runtime.Filters;
+using Elsa.Workflows.Runtime.Stores;
+using Microsoft.Extensions.Logging.Abstractions;
+using NSubstitute;
+
+namespace Elsa.Scheduling.UnitTests.Services;
+
+public class SchedulingBookmarkReconcilerTests
+{
+ [Fact]
+ public async Task ClassifyAsync_SkipsMissingBlankAndFinishedInstances()
+ {
+ var instanceStore = CreateInstanceStore(
+ Instance("running", WorkflowStatus.Running, WorkflowSubStatus.Executing),
+ Instance("suspended", WorkflowStatus.Running, WorkflowSubStatus.Suspended),
+ Instance("interrupted", WorkflowStatus.Running, WorkflowSubStatus.Interrupted),
+ Instance("pending", WorkflowStatus.Running, WorkflowSubStatus.Pending),
+ Instance("finished", WorkflowStatus.Finished, WorkflowSubStatus.Finished),
+ Instance("cancelled", WorkflowStatus.Finished, WorkflowSubStatus.Cancelled),
+ Instance("faulted", WorkflowStatus.Finished, WorkflowSubStatus.Faulted));
+ var reconciler = new SchedulingBookmarkReconciler(instanceStore, Substitute.For());
+
+ var result = await reconciler.ClassifyAsync(
+ [
+ Bookmark("b-missing", "missing"),
+ Bookmark("b-blank", ""),
+ Bookmark("b-finished", "finished"),
+ Bookmark("b-cancelled", "cancelled"),
+ Bookmark("b-faulted", "faulted"),
+ Bookmark("b-running", "running"),
+ Bookmark("b-suspended", "suspended"),
+ Bookmark("b-interrupted", "interrupted"),
+ Bookmark("b-pending", "pending")
+ ]);
+
+ Assert.Equal(["b-running", "b-suspended", "b-interrupted", "b-pending"], result.Schedulable.Select(x => x.Id));
+ Assert.Equal(["b-missing", "b-blank", "b-finished", "b-cancelled", "b-faulted"], result.Orphans.Select(x => x.Id));
+ }
+
+ [Fact]
+ public async Task ClassifyAsync_KeepsPastDueSuspendedDelayBookmarks()
+ {
+ var instanceStore = CreateInstanceStore(Instance("suspended", WorkflowStatus.Running, WorkflowSubStatus.Suspended));
+ var reconciler = new SchedulingBookmarkReconciler(instanceStore, Substitute.For());
+
+ var result = await reconciler.ClassifyAsync([Bookmark("past-due-delay", "suspended")]);
+
+ Assert.Equal(["past-due-delay"], result.Schedulable.Select(x => x.Id));
+ Assert.Empty(result.Orphans);
+ }
+
+ [Fact]
+ public async Task PurgeAsync_DeletesOrphanBookmarksFromTheStore()
+ {
+ var bookmarkStore = new MemoryBookmarkStore(new MemoryStore());
+ await bookmarkStore.SaveManyAsync(
+ [
+ Bookmark("orphan-1", "missing"),
+ Bookmark("orphan-2", "finished"),
+ Bookmark("live", "suspended")
+ ], CancellationToken.None);
+ var bookmarkManager = new DefaultBookmarkManager(bookmarkStore, Substitute.For(), NullLogger.Instance);
+ var reconciler = new SchedulingBookmarkReconciler(CreateInstanceStore(), bookmarkManager);
+
+ await reconciler.PurgeAsync([Bookmark("orphan-1", "missing"), Bookmark("orphan-2", "finished")]);
+
+ var remaining = (await bookmarkStore.FindManyAsync(new BookmarkFilter())).Select(x => x.Id).ToList();
+ Assert.Equal(["live"], remaining);
+ }
+
+ [Fact]
+ public async Task PurgeAsync_DoesNotDeleteWhenThereAreNoOrphans()
+ {
+ var bookmarkManager = Substitute.For();
+ var reconciler = new SchedulingBookmarkReconciler(CreateInstanceStore(), bookmarkManager);
+
+ await reconciler.PurgeAsync([]);
+
+ await bookmarkManager.DidNotReceive().DeleteManyAsync(Arg.Any(), Arg.Any());
+ }
+
+ [Fact]
+ public async Task PurgeAsync_DoesNotDeleteWhenInstanceBecameRunning()
+ {
+ var bookmarkManager = Substitute.For();
+ var reconciler = new SchedulingBookmarkReconciler(
+ CreateInstanceStore(Instance("revived", WorkflowStatus.Running, WorkflowSubStatus.Suspended)),
+ bookmarkManager);
+
+ await reconciler.PurgeAsync([Bookmark("revived-bookmark", "revived")]);
+
+ await bookmarkManager.DidNotReceive().DeleteManyAsync(Arg.Any(), Arg.Any());
+ }
+
+ [Fact]
+ public async Task PurgeAsync_DoesNotDeleteWhenBookmarkManagerIsMissing()
+ {
+ var reconciler = new SchedulingBookmarkReconciler(CreateInstanceStore(), bookmarkManager: null);
+
+ await reconciler.PurgeAsync([Bookmark("orphan-1", "missing")]);
+ }
+
+ private static MemoryWorkflowInstanceStore CreateInstanceStore(params WorkflowInstance[] instances)
+ {
+ var store = new MemoryStore();
+ store.AddMany(instances, x => x.Id);
+ return new MemoryWorkflowInstanceStore(store);
+ }
+
+ private static WorkflowInstance Instance(string id, WorkflowStatus status, WorkflowSubStatus subStatus) => new()
+ {
+ Id = id,
+ DefinitionId = "definition",
+ DefinitionVersionId = "version",
+ Status = status,
+ SubStatus = subStatus
+ };
+
+ private static StoredBookmark Bookmark(string id, string workflowInstanceId) => new()
+ {
+ Id = id,
+ Hash = "hash",
+ Name = SchedulingStimulusNames.Delay,
+ WorkflowInstanceId = workflowInstanceId
+ };
+}
diff --git a/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs b/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs
index cdeee7dbe..7390c1245 100644
--- a/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs
+++ b/test/unit/Elsa.Scheduling.UnitTests/StartupTasks/CreateSchedulesStartupTaskTests.cs
@@ -3,6 +3,8 @@ 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;
@@ -25,6 +27,8 @@ public class CreateSchedulesStartupTaskTests
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()
@@ -33,6 +37,8 @@ public class CreateSchedulesStartupTaskTests
.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]
@@ -101,14 +107,197 @@ public class CreateSchedulesStartupTaskTests
await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is>(x => x.SequenceEqual(secondBookmarkPage)), Arg.Any());
}
- private ServiceProvider CreateServiceProvider(Action? configureServices = null)
+ [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
+ };
}