Fix scheduling startup backlog catch-up
Rebuild local schedules in bounded pages and stagger past-due specific-instant catch-up so orphaned scheduling bookmarks do not flood dispatch during startup. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
parent
6a4c141ce1
commit
c7912fd2c8
|
|
@ -1,3 +1,5 @@
|
||||||
|
using Elsa.Common.Models;
|
||||||
|
using Elsa.Extensions;
|
||||||
using Elsa.Workflows;
|
using Elsa.Workflows;
|
||||||
using Elsa.Workflows.Runtime;
|
using Elsa.Workflows.Runtime;
|
||||||
using Elsa.Workflows.Runtime.Entities;
|
using Elsa.Workflows.Runtime.Entities;
|
||||||
|
|
@ -36,6 +38,14 @@ public class EFCoreBookmarkStore(Store<RuntimeElsaDbContext, StoredBookmark> sto
|
||||||
return await store.QueryAsync(filter.Apply, OnLoadAsync, filter.TenantAgnostic, cancellationToken);
|
return await store.QueryAsync(filter.Apply, OnLoadAsync, filter.TenantAgnostic, cancellationToken);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <inheritdoc />
|
||||||
|
public async ValueTask<Page<StoredBookmark>> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
var count = await store.CountAsync(filter.Apply, filter.TenantAgnostic, cancellationToken);
|
||||||
|
var results = (await store.QueryAsync(query => filter.Apply(query).OrderBy(x => x.Id).Paginate(pageArgs), OnLoadAsync, filter.TenantAgnostic, cancellationToken)).ToList();
|
||||||
|
return Page.Of(results, count);
|
||||||
|
}
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public async ValueTask<long> DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default)
|
public async ValueTask<long> DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -8,7 +8,6 @@ using Elsa.Workflows.Runtime.Entities;
|
||||||
using Elsa.Workflows.Runtime.Filters;
|
using Elsa.Workflows.Runtime.Filters;
|
||||||
using Elsa.Workflows.Runtime.OrderDefinitions;
|
using Elsa.Workflows.Runtime.OrderDefinitions;
|
||||||
using JetBrains.Annotations;
|
using JetBrains.Annotations;
|
||||||
using Open.Linq.AsyncExtensions;
|
|
||||||
|
|
||||||
namespace Elsa.Persistence.EFCore.Modules.Runtime;
|
namespace Elsa.Persistence.EFCore.Modules.Runtime;
|
||||||
|
|
||||||
|
|
@ -50,8 +49,8 @@ public class EFCoreTriggerStore(
|
||||||
|
|
||||||
public async ValueTask<Page<StoredTrigger>> FindManyAsync<TProp>(TriggerFilter filter, PageArgs pageArgs, StoredTriggerOrder<TProp> order, CancellationToken cancellationToken = default)
|
public async ValueTask<Page<StoredTrigger>> FindManyAsync<TProp>(TriggerFilter filter, PageArgs pageArgs, StoredTriggerOrder<TProp> order, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
var count = await store.QueryAsync(filter.Apply, OnLoadAsync, cancellationToken).LongCount();
|
var count = await store.CountAsync(filter.Apply, filter.TenantAgnostic, cancellationToken);
|
||||||
var results = await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(order).Paginate(pageArgs).OrderBy(order), OnLoadAsync, cancellationToken).ToList();
|
var results = (await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(order).Paginate(pageArgs).OrderBy(order), OnLoadAsync, filter.TenantAgnostic, cancellationToken)).ToList();
|
||||||
return new(results, count);
|
return new(results, count);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -6,6 +6,7 @@ using Elsa.Features.Attributes;
|
||||||
using Elsa.Features.Services;
|
using Elsa.Features.Services;
|
||||||
using Elsa.Scheduling.Bookmarks;
|
using Elsa.Scheduling.Bookmarks;
|
||||||
using Elsa.Scheduling.Handlers;
|
using Elsa.Scheduling.Handlers;
|
||||||
|
using Elsa.Scheduling.Options;
|
||||||
using Elsa.Scheduling.Services;
|
using Elsa.Scheduling.Services;
|
||||||
using Elsa.Scheduling.StartupTasks;
|
using Elsa.Scheduling.StartupTasks;
|
||||||
using Elsa.Scheduling.TriggerPayloadValidators;
|
using Elsa.Scheduling.TriggerPayloadValidators;
|
||||||
|
|
@ -44,6 +45,7 @@ public class SchedulingFeature : FeatureBase
|
||||||
.AddSingleton<UpdateTenantSchedules>()
|
.AddSingleton<UpdateTenantSchedules>()
|
||||||
.AddSingleton<ITenantDeletedEvent>(sp => sp.GetRequiredService<UpdateTenantSchedules>())
|
.AddSingleton<ITenantDeletedEvent>(sp => sp.GetRequiredService<UpdateTenantSchedules>())
|
||||||
.AddSingleton<IScheduler, LocalScheduler>()
|
.AddSingleton<IScheduler, LocalScheduler>()
|
||||||
|
.AddSingleton<PastDueScheduleStaggerer>()
|
||||||
.AddSingleton<CronosCronParser>()
|
.AddSingleton<CronosCronParser>()
|
||||||
.AddSingleton(CronParser)
|
.AddSingleton(CronParser)
|
||||||
.AddScoped<ITriggerScheduler, DefaultTriggerScheduler>()
|
.AddScoped<ITriggerScheduler, DefaultTriggerScheduler>()
|
||||||
|
|
@ -59,6 +61,7 @@ public class SchedulingFeature : FeatureBase
|
||||||
// Graceful shutdown: register scheduled-trigger ingress for diagnostic visibility (FR-006).
|
// Graceful shutdown: register scheduled-trigger ingress for diagnostic visibility (FR-006).
|
||||||
.AddSingleton<Elsa.Workflows.Runtime.IIngressSource, Elsa.Scheduling.IngressSources.ScheduledTriggerIngressSource>();
|
.AddSingleton<Elsa.Workflows.Runtime.IIngressSource, Elsa.Scheduling.IngressSources.ScheduledTriggerIngressSource>();
|
||||||
|
|
||||||
|
Services.Configure<SchedulingOptions>(_ => { });
|
||||||
Services.Configure<SerializationTypeOptions>(options =>
|
Services.Configure<SerializationTypeOptions>(options =>
|
||||||
{
|
{
|
||||||
options.RegisterTypeAlias(typeof(CronBookmarkPayload), nameof(CronBookmarkPayload));
|
options.RegisterTypeAlias(typeof(CronBookmarkPayload), nameof(CronBookmarkPayload));
|
||||||
|
|
|
||||||
27
src/modules/Elsa.Scheduling/Options/SchedulingOptions.cs
Normal file
27
src/modules/Elsa.Scheduling/Options/SchedulingOptions.cs
Normal file
|
|
@ -0,0 +1,27 @@
|
||||||
|
namespace Elsa.Scheduling.Options;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Configures local scheduling behavior.
|
||||||
|
/// </summary>
|
||||||
|
public class SchedulingOptions
|
||||||
|
{
|
||||||
|
/// <summary>
|
||||||
|
/// The number of stored triggers and bookmarks to load per batch when rebuilding local schedules on startup.
|
||||||
|
/// </summary>
|
||||||
|
public int StartupSchedulePageSize { get; set; } = 1000;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The minimum delay used when a specific-instant schedule is already due.
|
||||||
|
/// </summary>
|
||||||
|
public TimeSpan MinimumPastDueScheduleDelay { get; set; } = TimeSpan.FromMilliseconds(1);
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The spacing between past-due schedules during catch-up.
|
||||||
|
/// </summary>
|
||||||
|
public TimeSpan PastDueScheduleStaggerInterval { get; set; } = TimeSpan.FromMilliseconds(50);
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The bounded window over which past-due schedules are distributed.
|
||||||
|
/// </summary>
|
||||||
|
public TimeSpan PastDueScheduleStaggerWindow { get; set; } = TimeSpan.FromMinutes(5);
|
||||||
|
}
|
||||||
|
|
@ -1,9 +1,12 @@
|
||||||
using Elsa.Common;
|
using Elsa.Common;
|
||||||
using Elsa.Mediator.Contracts;
|
using Elsa.Mediator.Contracts;
|
||||||
using Elsa.Scheduling.Commands;
|
using Elsa.Scheduling.Commands;
|
||||||
|
using Elsa.Scheduling.Options;
|
||||||
|
using Elsa.Scheduling.Services;
|
||||||
using Microsoft.Extensions.DependencyInjection;
|
using Microsoft.Extensions.DependencyInjection;
|
||||||
using Microsoft.Extensions.Logging;
|
using Microsoft.Extensions.Logging;
|
||||||
using Timer = System.Timers.Timer;
|
using Timer = System.Timers.Timer;
|
||||||
|
using OptionsFactory = Microsoft.Extensions.Options.Options;
|
||||||
|
|
||||||
namespace Elsa.Scheduling.ScheduledTasks;
|
namespace Elsa.Scheduling.ScheduledTasks;
|
||||||
|
|
||||||
|
|
@ -12,10 +15,12 @@ namespace Elsa.Scheduling.ScheduledTasks;
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable
|
public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable
|
||||||
{
|
{
|
||||||
|
private static readonly PastDueScheduleStaggerer DefaultPastDueScheduleStaggerer = new(OptionsFactory.Create(new SchedulingOptions()));
|
||||||
private readonly ITask _task;
|
private readonly ITask _task;
|
||||||
private readonly ISystemClock _systemClock;
|
private readonly ISystemClock _systemClock;
|
||||||
private readonly IServiceScopeFactory _scopeFactory;
|
private readonly IServiceScopeFactory _scopeFactory;
|
||||||
private readonly ILogger<ScheduledSpecificInstantTask> _logger;
|
private readonly ILogger<ScheduledSpecificInstantTask> _logger;
|
||||||
|
private readonly PastDueScheduleStaggerer _pastDueScheduleStaggerer;
|
||||||
private readonly DateTimeOffset _startAt;
|
private readonly DateTimeOffset _startAt;
|
||||||
private readonly CancellationTokenSource _cancellationTokenSource;
|
private readonly CancellationTokenSource _cancellationTokenSource;
|
||||||
private readonly SemaphoreSlim _executionSemaphore = new(1, 1);
|
private readonly SemaphoreSlim _executionSemaphore = new(1, 1);
|
||||||
|
|
@ -27,12 +32,33 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Initializes a new instance of <see cref="ScheduledSpecificInstantTask"/>.
|
/// Initializes a new instance of <see cref="ScheduledSpecificInstantTask"/>.
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public ScheduledSpecificInstantTask(ITask task, DateTimeOffset startAt, ISystemClock systemClock, IServiceScopeFactory scopeFactory, ILogger<ScheduledSpecificInstantTask> logger)
|
public ScheduledSpecificInstantTask(
|
||||||
|
ITask task,
|
||||||
|
DateTimeOffset startAt,
|
||||||
|
ISystemClock systemClock,
|
||||||
|
IServiceScopeFactory scopeFactory,
|
||||||
|
ILogger<ScheduledSpecificInstantTask> logger)
|
||||||
|
: this(task, startAt, systemClock, scopeFactory, logger, DefaultPastDueScheduleStaggerer)
|
||||||
|
{
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Initializes a new instance of <see cref="ScheduledSpecificInstantTask"/>.
|
||||||
|
/// </summary>
|
||||||
|
[ActivatorUtilitiesConstructor]
|
||||||
|
public ScheduledSpecificInstantTask(
|
||||||
|
ITask task,
|
||||||
|
DateTimeOffset startAt,
|
||||||
|
ISystemClock systemClock,
|
||||||
|
IServiceScopeFactory scopeFactory,
|
||||||
|
ILogger<ScheduledSpecificInstantTask> logger,
|
||||||
|
PastDueScheduleStaggerer pastDueScheduleStaggerer)
|
||||||
{
|
{
|
||||||
_task = task;
|
_task = task;
|
||||||
_systemClock = systemClock;
|
_systemClock = systemClock;
|
||||||
_scopeFactory = scopeFactory;
|
_scopeFactory = scopeFactory;
|
||||||
_logger = logger;
|
_logger = logger;
|
||||||
|
_pastDueScheduleStaggerer = pastDueScheduleStaggerer;
|
||||||
_startAt = startAt;
|
_startAt = startAt;
|
||||||
_cancellationTokenSource = new();
|
_cancellationTokenSource = new();
|
||||||
|
|
||||||
|
|
@ -57,16 +83,14 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable
|
||||||
{
|
{
|
||||||
var now = _systemClock.UtcNow;
|
var now = _systemClock.UtcNow;
|
||||||
var delay = _startAt - now;
|
var delay = _startAt - now;
|
||||||
|
var adjustedDelay = _pastDueScheduleStaggerer.GetDelay(delay);
|
||||||
|
|
||||||
// Handle edge cases where delay is zero or negative (e.g., due to clock drift, fast execution, or time alignment)
|
|
||||||
// Instead of silently returning, use a minimum delay to ensure the timer fires and workflow continues scheduling
|
|
||||||
if (delay <= TimeSpan.Zero)
|
if (delay <= TimeSpan.Zero)
|
||||||
{
|
{
|
||||||
_logger.LogWarning("Calculated delay is {Delay} which is not positive. Using minimum delay of 1ms to ensure timer fires", delay);
|
_logger.LogDebug("Calculated delay is {Delay} which is not positive. Using catch-up delay of {CatchUpDelay}", delay, adjustedDelay);
|
||||||
delay = TimeSpan.FromMilliseconds(1);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
_timer = new(delay.TotalMilliseconds)
|
_timer = new(adjustedDelay.TotalMilliseconds)
|
||||||
{
|
{
|
||||||
Enabled = true
|
Enabled = true
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,34 @@
|
||||||
|
using Elsa.Scheduling.Options;
|
||||||
|
using Microsoft.Extensions.Options;
|
||||||
|
|
||||||
|
namespace Elsa.Scheduling.Services;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Distributes already-due schedules over a bounded window to avoid dispatch storms during startup catch-up.
|
||||||
|
/// </summary>
|
||||||
|
public class PastDueScheduleStaggerer(IOptions<SchedulingOptions> options)
|
||||||
|
{
|
||||||
|
private long _sequence;
|
||||||
|
|
||||||
|
public TimeSpan GetDelay(TimeSpan calculatedDelay)
|
||||||
|
{
|
||||||
|
if (calculatedDelay > TimeSpan.Zero)
|
||||||
|
return calculatedDelay;
|
||||||
|
|
||||||
|
var currentOptions = options.Value;
|
||||||
|
var minimumDelay = GetPositiveOrDefault(currentOptions.MinimumPastDueScheduleDelay, TimeSpan.FromMilliseconds(1));
|
||||||
|
var staggerInterval = currentOptions.PastDueScheduleStaggerInterval;
|
||||||
|
var staggerWindow = currentOptions.PastDueScheduleStaggerWindow;
|
||||||
|
|
||||||
|
if (staggerInterval <= TimeSpan.Zero || staggerWindow <= TimeSpan.Zero)
|
||||||
|
return minimumDelay;
|
||||||
|
|
||||||
|
var slotCount = Math.Max(1, staggerWindow.Ticks / staggerInterval.Ticks);
|
||||||
|
var sequence = Interlocked.Increment(ref _sequence) - 1;
|
||||||
|
var slot = (sequence & long.MaxValue) % slotCount;
|
||||||
|
|
||||||
|
return minimumDelay + TimeSpan.FromTicks(staggerInterval.Ticks * slot);
|
||||||
|
}
|
||||||
|
|
||||||
|
private static TimeSpan GetPositiveOrDefault(TimeSpan value, TimeSpan defaultValue) => value > TimeSpan.Zero ? value : defaultValue;
|
||||||
|
}
|
||||||
|
|
@ -5,6 +5,7 @@ using Elsa.Extensions;
|
||||||
using Elsa.Workflows.Options;
|
using Elsa.Workflows.Options;
|
||||||
using Elsa.Scheduling.Bookmarks;
|
using Elsa.Scheduling.Bookmarks;
|
||||||
using Elsa.Scheduling.Handlers;
|
using Elsa.Scheduling.Handlers;
|
||||||
|
using Elsa.Scheduling.Options;
|
||||||
using Elsa.Scheduling.Services;
|
using Elsa.Scheduling.Services;
|
||||||
using Elsa.Scheduling.StartupTasks;
|
using Elsa.Scheduling.StartupTasks;
|
||||||
using Elsa.Scheduling.TriggerPayloadValidators;
|
using Elsa.Scheduling.TriggerPayloadValidators;
|
||||||
|
|
@ -44,6 +45,7 @@ public class SchedulingFeature : IShellFeature
|
||||||
.AddSingleton<UpdateTenantSchedules>()
|
.AddSingleton<UpdateTenantSchedules>()
|
||||||
.AddSingleton<ITenantDeletedEvent>(sp => sp.GetRequiredService<UpdateTenantSchedules>())
|
.AddSingleton<ITenantDeletedEvent>(sp => sp.GetRequiredService<UpdateTenantSchedules>())
|
||||||
.AddSingleton<IScheduler, LocalScheduler>()
|
.AddSingleton<IScheduler, LocalScheduler>()
|
||||||
|
.AddSingleton<PastDueScheduleStaggerer>()
|
||||||
.AddSingleton<CronosCronParser>()
|
.AddSingleton<CronosCronParser>()
|
||||||
.AddSingleton(CronParser)
|
.AddSingleton(CronParser)
|
||||||
.AddScoped<ITriggerScheduler, DefaultTriggerScheduler>()
|
.AddScoped<ITriggerScheduler, DefaultTriggerScheduler>()
|
||||||
|
|
@ -55,6 +57,7 @@ public class SchedulingFeature : IShellFeature
|
||||||
.AddTriggerPayloadValidator<CronTriggerPayloadValidator, CronTriggerPayload>()
|
.AddTriggerPayloadValidator<CronTriggerPayloadValidator, CronTriggerPayload>()
|
||||||
.AddActivitiesFrom<SchedulingFeature>();
|
.AddActivitiesFrom<SchedulingFeature>();
|
||||||
|
|
||||||
|
services.Configure<SchedulingOptions>(_ => { });
|
||||||
services.Configure<SerializationTypeOptions>(options =>
|
services.Configure<SerializationTypeOptions>(options =>
|
||||||
{
|
{
|
||||||
options.AddTypeAlias<CronBookmarkPayload>();
|
options.AddTypeAlias<CronBookmarkPayload>();
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,12 @@
|
||||||
using Elsa.Common;
|
using Elsa.Common;
|
||||||
using Elsa.Common.Multitenancy;
|
using Elsa.Common.Multitenancy;
|
||||||
|
using Elsa.Common.Models;
|
||||||
|
using Elsa.Scheduling.Options;
|
||||||
using Elsa.Workflows.Runtime;
|
using Elsa.Workflows.Runtime;
|
||||||
using Elsa.Workflows.Runtime.Filters;
|
using Elsa.Workflows.Runtime.Filters;
|
||||||
using Elsa.Workflows.Runtime.Tasks;
|
using Elsa.Workflows.Runtime.Tasks;
|
||||||
using Microsoft.Extensions.DependencyInjection;
|
using Microsoft.Extensions.DependencyInjection;
|
||||||
|
using Microsoft.Extensions.Options;
|
||||||
|
|
||||||
namespace Elsa.Scheduling.StartupTasks;
|
namespace Elsa.Scheduling.StartupTasks;
|
||||||
|
|
||||||
|
|
@ -11,7 +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.
|
/// Enqueues schedule creation when using the default scheduler, which doesn't have its own persistence layer like Quartz or Hangfire.
|
||||||
/// </summary>
|
/// </summary>
|
||||||
[TaskDependency(typeof(PopulateRegistriesStartupTask))]
|
[TaskDependency(typeof(PopulateRegistriesStartupTask))]
|
||||||
public class CreateSchedulesStartupTask(IServiceProvider serviceProvider) : IStartupTask
|
public class CreateSchedulesStartupTask(IServiceProvider serviceProvider, IOptions<SchedulingOptions> options) : IStartupTask
|
||||||
{
|
{
|
||||||
public async Task ExecuteAsync(CancellationToken cancellationToken)
|
public async Task ExecuteAsync(CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
|
|
@ -23,12 +26,13 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider) : ISta
|
||||||
await CreateSchedulesAsync(serviceProvider, cancellationToken);
|
await CreateSchedulesAsync(serviceProvider, cancellationToken);
|
||||||
}
|
}
|
||||||
|
|
||||||
private static async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
|
private async Task CreateSchedulesAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
var triggerStore = serviceProvider.GetRequiredService<ITriggerStore>();
|
var triggerStore = serviceProvider.GetRequiredService<ITriggerStore>();
|
||||||
var bookmarkStore = serviceProvider.GetRequiredService<IBookmarkStore>();
|
var bookmarkStore = serviceProvider.GetRequiredService<IBookmarkStore>();
|
||||||
var triggerScheduler = serviceProvider.GetRequiredService<ITriggerScheduler>();
|
var triggerScheduler = serviceProvider.GetRequiredService<ITriggerScheduler>();
|
||||||
var bookmarkScheduler = serviceProvider.GetRequiredService<IBookmarkScheduler>();
|
var bookmarkScheduler = serviceProvider.GetRequiredService<IBookmarkScheduler>();
|
||||||
|
var pageSize = Math.Max(1, options.Value.StartupSchedulePageSize);
|
||||||
var stimulusNames = new[]
|
var stimulusNames = new[]
|
||||||
{
|
{
|
||||||
SchedulingStimulusNames.Cron, SchedulingStimulusNames.Timer, SchedulingStimulusNames.StartAt, SchedulingStimulusNames.Delay,
|
SchedulingStimulusNames.Cron, SchedulingStimulusNames.Timer, SchedulingStimulusNames.StartAt, SchedulingStimulusNames.Delay,
|
||||||
|
|
@ -41,10 +45,50 @@ public class CreateSchedulesStartupTask(IServiceProvider serviceProvider) : ISta
|
||||||
{
|
{
|
||||||
Names = stimulusNames
|
Names = stimulusNames
|
||||||
};
|
};
|
||||||
var triggers = (await triggerStore.FindManyAsync(triggerFilter, cancellationToken)).ToList();
|
|
||||||
var bookmarks = (await bookmarkStore.FindManyAsync(bookmarkFilter, cancellationToken)).ToList();
|
|
||||||
|
|
||||||
await triggerScheduler.ScheduleAsync(triggers, cancellationToken);
|
await ScheduleTriggersAsync(triggerStore, triggerScheduler, triggerFilter, pageSize, cancellationToken);
|
||||||
await bookmarkScheduler.ScheduleAsync(bookmarks, cancellationToken);
|
await ScheduleBookmarksAsync(bookmarkStore, bookmarkScheduler, 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, BookmarkFilter bookmarkFilter, int pageSize, CancellationToken cancellationToken)
|
||||||
|
{
|
||||||
|
var pageArgs = PageArgs.FromRange(0, pageSize);
|
||||||
|
|
||||||
|
while (true)
|
||||||
|
{
|
||||||
|
var page = await bookmarkStore.FindManyAsync(bookmarkFilter, pageArgs, cancellationToken);
|
||||||
|
|
||||||
|
if (page.Items.Count == 0)
|
||||||
|
break;
|
||||||
|
|
||||||
|
await bookmarkScheduler.ScheduleAsync(page.Items, cancellationToken);
|
||||||
|
|
||||||
|
var nextOffset = pageArgs.Offset.GetValueOrDefault() + page.Items.Count;
|
||||||
|
if (nextOffset >= page.TotalCount)
|
||||||
|
break;
|
||||||
|
|
||||||
|
pageArgs = pageArgs.Next();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,3 +1,4 @@
|
||||||
|
using Elsa.Common.Models;
|
||||||
using Elsa.Workflows.Runtime.Entities;
|
using Elsa.Workflows.Runtime.Entities;
|
||||||
using Elsa.Workflows.Runtime.Filters;
|
using Elsa.Workflows.Runtime.Filters;
|
||||||
|
|
||||||
|
|
@ -34,6 +35,26 @@ public interface IBookmarkStore
|
||||||
/// </summary>
|
/// </summary>
|
||||||
ValueTask<IEnumerable<StoredBookmark>> FindManyAsync(BookmarkFilter filter, CancellationToken cancellationToken = default);
|
ValueTask<IEnumerable<StoredBookmark>> FindManyAsync(BookmarkFilter filter, CancellationToken cancellationToken = default);
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Returns a page of bookmarks matching the specified filter.
|
||||||
|
/// </summary>
|
||||||
|
/// <remarks>
|
||||||
|
/// The default implementation materializes all matching bookmarks and pages them in memory. Stores backed by external persistence should override this method.
|
||||||
|
/// </remarks>
|
||||||
|
async ValueTask<Page<StoredBookmark>> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
var records = (await FindManyAsync(filter, cancellationToken)).OrderBy(x => x.Id).ToList();
|
||||||
|
IEnumerable<StoredBookmark> page = records;
|
||||||
|
|
||||||
|
if (pageArgs.Offset.HasValue)
|
||||||
|
page = page.Skip(pageArgs.Offset.Value);
|
||||||
|
|
||||||
|
if (pageArgs.Limit.HasValue)
|
||||||
|
page = page.Take(pageArgs.Limit.Value);
|
||||||
|
|
||||||
|
return Page.Of(page.ToList(), records.Count);
|
||||||
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Deletes a set of bookmarks matching the specified filter.
|
/// Deletes a set of bookmarks matching the specified filter.
|
||||||
/// </summary>
|
/// </summary>
|
||||||
|
|
|
||||||
|
|
@ -48,13 +48,13 @@ public class CachingTriggerStore(ITriggerStore decoratedStore, ICacheManager cac
|
||||||
|
|
||||||
public async ValueTask<Page<StoredTrigger>> FindManyAsync(TriggerFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default)
|
public async ValueTask<Page<StoredTrigger>> FindManyAsync(TriggerFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
var cacheKey = hasher.Hash(filter);
|
var cacheKey = hasher.Hash(filter, pageArgs);
|
||||||
return (await GetOrCreateAsync(cacheKey, async () => await decoratedStore.FindManyAsync(filter, pageArgs, cancellationToken)))!;
|
return (await GetOrCreateAsync(cacheKey, async () => await decoratedStore.FindManyAsync(filter, pageArgs, cancellationToken)))!;
|
||||||
}
|
}
|
||||||
|
|
||||||
public async ValueTask<Page<StoredTrigger>> FindManyAsync<TProp>(TriggerFilter filter, PageArgs pageArgs, StoredTriggerOrder<TProp> order, CancellationToken cancellationToken = default)
|
public async ValueTask<Page<StoredTrigger>> FindManyAsync<TProp>(TriggerFilter filter, PageArgs pageArgs, StoredTriggerOrder<TProp> order, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
var cacheKey = hasher.Hash(filter);
|
var cacheKey = hasher.Hash(filter, pageArgs, order);
|
||||||
return (await GetOrCreateAsync(cacheKey, async () => await decoratedStore.FindManyAsync(filter, pageArgs, order, cancellationToken)))!;
|
return (await GetOrCreateAsync(cacheKey, async () => await decoratedStore.FindManyAsync(filter, pageArgs, order, cancellationToken)))!;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,6 @@
|
||||||
|
using Elsa.Common.Models;
|
||||||
using Elsa.Common.Services;
|
using Elsa.Common.Services;
|
||||||
|
using Elsa.Extensions;
|
||||||
using Elsa.Workflows.Runtime.Entities;
|
using Elsa.Workflows.Runtime.Entities;
|
||||||
using Elsa.Workflows.Runtime.Filters;
|
using Elsa.Workflows.Runtime.Filters;
|
||||||
using JetBrains.Annotations;
|
using JetBrains.Annotations;
|
||||||
|
|
@ -37,6 +39,14 @@ public class MemoryBookmarkStore(MemoryStore<StoredBookmark> store) : IBookmarkS
|
||||||
return new(entities);
|
return new(entities);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <inheritdoc />
|
||||||
|
public ValueTask<Page<StoredBookmark>> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
var count = store.Query(query => Filter(query, filter)).LongCount();
|
||||||
|
var result = store.Query(query => Filter(query, filter).OrderBy(x => x.Id).Paginate(pageArgs)).ToList();
|
||||||
|
return new(Page.Of(result, count));
|
||||||
|
}
|
||||||
|
|
||||||
/// <inheritdoc />
|
/// <inheritdoc />
|
||||||
public async ValueTask<long> DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default)
|
public async ValueTask<long> DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
using System.Collections;
|
using System.Collections;
|
||||||
|
using Elsa.Common.Models;
|
||||||
using Elsa.Http.Bookmarks;
|
using Elsa.Http.Bookmarks;
|
||||||
using Elsa.Http.Middleware;
|
using Elsa.Http.Middleware;
|
||||||
using Elsa.Http.Options;
|
using Elsa.Http.Options;
|
||||||
|
|
@ -104,6 +105,13 @@ public class HttpWorkflowsMiddlewareTests
|
||||||
return new([]);
|
return new([]);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public ValueTask<Page<StoredBookmark>> FindManyAsync(BookmarkFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
LastFilter = filter;
|
||||||
|
var results = Filter(filter).ToList();
|
||||||
|
return new(Page.Of(results, results.Count));
|
||||||
|
}
|
||||||
|
|
||||||
public ValueTask<long> DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default)
|
public ValueTask<long> DeleteAsync(BookmarkFilter filter, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
var bookmarksToDelete = Filter(filter).ToList();
|
var bookmarksToDelete = Filter(filter).ToList();
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,12 @@
|
||||||
using Elsa.Common;
|
using Elsa.Common;
|
||||||
using Elsa.Mediator.Contracts;
|
using Elsa.Mediator.Contracts;
|
||||||
|
using Elsa.Scheduling.Options;
|
||||||
using Elsa.Scheduling.ScheduledTasks;
|
using Elsa.Scheduling.ScheduledTasks;
|
||||||
|
using Elsa.Scheduling.Services;
|
||||||
using Microsoft.Extensions.DependencyInjection;
|
using Microsoft.Extensions.DependencyInjection;
|
||||||
using Microsoft.Extensions.Logging;
|
using Microsoft.Extensions.Logging;
|
||||||
using NSubstitute;
|
using NSubstitute;
|
||||||
|
using OptionsFactory = Microsoft.Extensions.Options.Options;
|
||||||
|
|
||||||
namespace Elsa.Scheduling.UnitTests.ScheduledTasks;
|
namespace Elsa.Scheduling.UnitTests.ScheduledTasks;
|
||||||
|
|
||||||
|
|
@ -17,6 +20,12 @@ public class ScheduledSpecificInstantTaskTests : IDisposable
|
||||||
private readonly ServiceProvider _serviceProvider;
|
private readonly ServiceProvider _serviceProvider;
|
||||||
private readonly ISystemClock _systemClock;
|
private readonly ISystemClock _systemClock;
|
||||||
private readonly ILogger<ScheduledSpecificInstantTask> _logger;
|
private readonly ILogger<ScheduledSpecificInstantTask> _logger;
|
||||||
|
private readonly PastDueScheduleStaggerer _pastDueScheduleStaggerer = new(OptionsFactory.Create(new SchedulingOptions
|
||||||
|
{
|
||||||
|
MinimumPastDueScheduleDelay = TimeSpan.FromMilliseconds(10),
|
||||||
|
PastDueScheduleStaggerInterval = TimeSpan.FromMilliseconds(25),
|
||||||
|
PastDueScheduleStaggerWindow = TimeSpan.FromMilliseconds(100)
|
||||||
|
}));
|
||||||
private readonly List<ScheduledSpecificInstantTask> _tasksToDispose = new();
|
private readonly List<ScheduledSpecificInstantTask> _tasksToDispose = new();
|
||||||
|
|
||||||
public ScheduledSpecificInstantTaskTests()
|
public ScheduledSpecificInstantTaskTests()
|
||||||
|
|
@ -39,7 +48,8 @@ public class ScheduledSpecificInstantTaskTests : IDisposable
|
||||||
startAt ?? DefaultNow.AddMinutes(5),
|
startAt ?? DefaultNow.AddMinutes(5),
|
||||||
systemClock ?? _systemClock,
|
systemClock ?? _systemClock,
|
||||||
_serviceProvider.CreateScope().ServiceProvider.GetRequiredService<IServiceScopeFactory>(),
|
_serviceProvider.CreateScope().ServiceProvider.GetRequiredService<IServiceScopeFactory>(),
|
||||||
_logger
|
_logger,
|
||||||
|
_pastDueScheduleStaggerer
|
||||||
);
|
);
|
||||||
_tasksToDispose.Add(scheduledTask);
|
_tasksToDispose.Add(scheduledTask);
|
||||||
return scheduledTask;
|
return scheduledTask;
|
||||||
|
|
@ -60,10 +70,10 @@ public class ScheduledSpecificInstantTaskTests : IDisposable
|
||||||
Arg.Any<Func<object, Exception?, string>>());
|
Arg.Any<Func<object, Exception?, string>>());
|
||||||
}
|
}
|
||||||
|
|
||||||
private void AssertWarningLogged(int expectedCount = 1)
|
private void AssertDebugLogged(int expectedCount = 1)
|
||||||
{
|
{
|
||||||
_logger.Received(expectedCount).Log(
|
_logger.Received(expectedCount).Log(
|
||||||
LogLevel.Warning,
|
LogLevel.Debug,
|
||||||
Arg.Any<EventId>(),
|
Arg.Any<EventId>(),
|
||||||
Arg.Any<object>(),
|
Arg.Any<object>(),
|
||||||
Arg.Any<Exception>(),
|
Arg.Any<Exception>(),
|
||||||
|
|
@ -85,31 +95,58 @@ public class ScheduledSpecificInstantTaskTests : IDisposable
|
||||||
}
|
}
|
||||||
|
|
||||||
[Fact]
|
[Fact]
|
||||||
public void Schedule_WithZeroDelay_ShouldUseMinimumDelay()
|
public void Schedule_WithZeroDelay_ShouldUseCatchUpDelay()
|
||||||
{
|
{
|
||||||
// Arrange - simulate a case where startAt is exactly now
|
// Arrange - simulate a case where startAt is exactly now
|
||||||
SetupSystemClock(DefaultNow);
|
SetupSystemClock(DefaultNow);
|
||||||
var startAt = DefaultNow; // delay = 0
|
var startAt = DefaultNow; // delay = 0
|
||||||
|
|
||||||
// Act - Should adjust to 1ms minimum delay
|
// Act - Should adjust to a bounded catch-up delay
|
||||||
CreateScheduledTask(startAt: startAt);
|
CreateScheduledTask(startAt: startAt);
|
||||||
|
|
||||||
// Assert - Should not crash and should log warning
|
AssertDebugLogged();
|
||||||
AssertWarningLogged();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
[Fact]
|
[Fact]
|
||||||
public void Schedule_WithNegativeDelay_ShouldUseMinimumDelay()
|
public void Schedule_WithNegativeDelay_ShouldUseCatchUpDelay()
|
||||||
{
|
{
|
||||||
// Arrange - simulate a case where startAt is in the past
|
// Arrange - simulate a case where startAt is in the past
|
||||||
SetupSystemClock(DefaultNow);
|
SetupSystemClock(DefaultNow);
|
||||||
var startAt = DefaultNow.AddMinutes(-1); // Past time
|
var startAt = DefaultNow.AddMinutes(-1); // Past time
|
||||||
|
|
||||||
// Act - Should adjust to 1ms minimum delay
|
// Act - Should adjust to a bounded catch-up delay
|
||||||
CreateScheduledTask(startAt: startAt);
|
CreateScheduledTask(startAt: startAt);
|
||||||
|
|
||||||
// Assert - Should log a warning
|
AssertDebugLogged();
|
||||||
AssertWarningLogged();
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public void Schedule_WithMultiplePastDueTasks_ShouldStaggerCatchUp()
|
||||||
|
{
|
||||||
|
SetupSystemClock(DefaultNow);
|
||||||
|
|
||||||
|
CreateScheduledTask(startAt: DefaultNow.AddMinutes(-1));
|
||||||
|
CreateScheduledTask(startAt: DefaultNow.AddMinutes(-1));
|
||||||
|
CreateScheduledTask(startAt: DefaultNow.AddMinutes(-1));
|
||||||
|
|
||||||
|
_logger.Received(1).Log(
|
||||||
|
LogLevel.Debug,
|
||||||
|
Arg.Any<EventId>(),
|
||||||
|
Arg.Is<object>(x => x.ToString()!.Contains("00:00:00.0100000")),
|
||||||
|
Arg.Any<Exception>(),
|
||||||
|
Arg.Any<Func<object, Exception?, string>>());
|
||||||
|
_logger.Received(1).Log(
|
||||||
|
LogLevel.Debug,
|
||||||
|
Arg.Any<EventId>(),
|
||||||
|
Arg.Is<object>(x => x.ToString()!.Contains("00:00:00.0350000")),
|
||||||
|
Arg.Any<Exception>(),
|
||||||
|
Arg.Any<Func<object, Exception?, string>>());
|
||||||
|
_logger.Received(1).Log(
|
||||||
|
LogLevel.Debug,
|
||||||
|
Arg.Any<EventId>(),
|
||||||
|
Arg.Is<object>(x => x.ToString()!.Contains("00:00:00.0600000")),
|
||||||
|
Arg.Any<Exception>(),
|
||||||
|
Arg.Any<Func<object, Exception?, string>>());
|
||||||
}
|
}
|
||||||
|
|
||||||
[Fact]
|
[Fact]
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,7 @@
|
||||||
using Elsa.Common;
|
using Elsa.Common;
|
||||||
using Elsa.Common.Multitenancy;
|
using Elsa.Common.Multitenancy;
|
||||||
|
using Elsa.Common.Models;
|
||||||
|
using Elsa.Scheduling.Options;
|
||||||
using Elsa.Scheduling.StartupTasks;
|
using Elsa.Scheduling.StartupTasks;
|
||||||
using Elsa.Workflows.Runtime;
|
using Elsa.Workflows.Runtime;
|
||||||
using Elsa.Workflows.Runtime.Entities;
|
using Elsa.Workflows.Runtime.Entities;
|
||||||
|
|
@ -7,22 +9,30 @@ using Elsa.Workflows.Runtime.Filters;
|
||||||
using Elsa.Workflows.Runtime.Tasks;
|
using Elsa.Workflows.Runtime.Tasks;
|
||||||
using Microsoft.Extensions.DependencyInjection;
|
using Microsoft.Extensions.DependencyInjection;
|
||||||
using NSubstitute;
|
using NSubstitute;
|
||||||
|
using OptionsFactory = Microsoft.Extensions.Options.Options;
|
||||||
|
|
||||||
namespace Elsa.Scheduling.UnitTests.StartupTasks;
|
namespace Elsa.Scheduling.UnitTests.StartupTasks;
|
||||||
|
|
||||||
public class CreateSchedulesStartupTaskTests
|
public class CreateSchedulesStartupTaskTests
|
||||||
{
|
{
|
||||||
private readonly StoredTrigger[] _triggers = [new() { WorkflowDefinitionId = "definition", WorkflowDefinitionVersionId = "version", ActivityId = "activity" }];
|
private readonly StoredTrigger[] _triggers =
|
||||||
private readonly StoredBookmark[] _bookmarks = [new() { Hash = "hash", WorkflowInstanceId = "instance" }];
|
[
|
||||||
|
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<ITriggerStore>();
|
private readonly ITriggerStore _triggerStore = Substitute.For<ITriggerStore>();
|
||||||
private readonly IBookmarkStore _bookmarkStore = Substitute.For<IBookmarkStore>();
|
private readonly IBookmarkStore _bookmarkStore = Substitute.For<IBookmarkStore>();
|
||||||
private readonly ITriggerScheduler _triggerScheduler = Substitute.For<ITriggerScheduler>();
|
private readonly ITriggerScheduler _triggerScheduler = Substitute.For<ITriggerScheduler>();
|
||||||
private readonly IBookmarkScheduler _bookmarkScheduler = Substitute.For<IBookmarkScheduler>();
|
private readonly IBookmarkScheduler _bookmarkScheduler = Substitute.For<IBookmarkScheduler>();
|
||||||
|
private readonly SchedulingOptions _options = new() { StartupSchedulePageSize = 1000 };
|
||||||
|
|
||||||
public CreateSchedulesStartupTaskTests()
|
public CreateSchedulesStartupTaskTests()
|
||||||
{
|
{
|
||||||
_triggerStore.FindManyAsync(Arg.Any<TriggerFilter>(), Arg.Any<CancellationToken>()).Returns(_triggers);
|
_triggerStore.FindManyAsync(Arg.Any<TriggerFilter>(), Arg.Any<PageArgs>(), Arg.Any<CancellationToken>())
|
||||||
_bookmarkStore.FindManyAsync(Arg.Any<BookmarkFilter>(), Arg.Any<CancellationToken>()).Returns(_bookmarks);
|
.Returns(new Page<StoredTrigger>(_triggers, _triggers.Length));
|
||||||
|
_bookmarkStore.FindManyAsync(Arg.Any<BookmarkFilter>(), Arg.Any<PageArgs>(), Arg.Any<CancellationToken>())
|
||||||
|
.Returns(new Page<StoredBookmark>(_bookmarks, _bookmarks.Length));
|
||||||
}
|
}
|
||||||
|
|
||||||
[Fact]
|
[Fact]
|
||||||
|
|
@ -36,7 +46,7 @@ public class CreateSchedulesStartupTaskTests
|
||||||
[Fact]
|
[Fact]
|
||||||
public async Task ExecuteAsync_WithoutTenantBackgroundQueue_SchedulesImmediately()
|
public async Task ExecuteAsync_WithoutTenantBackgroundQueue_SchedulesImmediately()
|
||||||
{
|
{
|
||||||
var task = new CreateSchedulesStartupTask(CreateServiceProvider());
|
var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options));
|
||||||
|
|
||||||
await task.ExecuteAsync(CancellationToken.None);
|
await task.ExecuteAsync(CancellationToken.None);
|
||||||
|
|
||||||
|
|
@ -51,7 +61,7 @@ public class CreateSchedulesStartupTaskTests
|
||||||
var workQueue = Substitute.For<ITenantBackgroundWorkQueue>();
|
var workQueue = Substitute.For<ITenantBackgroundWorkQueue>();
|
||||||
workQueue.EnqueueAsync(Arg.Do<TenantBackgroundWorkItem>(x => workItem = x), Arg.Any<CancellationToken>()).Returns(ValueTask.CompletedTask);
|
workQueue.EnqueueAsync(Arg.Do<TenantBackgroundWorkItem>(x => workItem = x), Arg.Any<CancellationToken>()).Returns(ValueTask.CompletedTask);
|
||||||
var serviceProvider = CreateServiceProvider(services => services.AddSingleton(workQueue));
|
var serviceProvider = CreateServiceProvider(services => services.AddSingleton(workQueue));
|
||||||
var task = new CreateSchedulesStartupTask(serviceProvider);
|
var task = new CreateSchedulesStartupTask(serviceProvider, OptionsFactory.Create(_options));
|
||||||
|
|
||||||
await task.ExecuteAsync(CancellationToken.None);
|
await task.ExecuteAsync(CancellationToken.None);
|
||||||
|
|
||||||
|
|
@ -65,6 +75,32 @@ public class CreateSchedulesStartupTaskTests
|
||||||
await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is<IEnumerable<StoredBookmark>>(x => x.SequenceEqual(_bookmarks)), Arg.Any<CancellationToken>());
|
await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is<IEnumerable<StoredBookmark>>(x => x.SequenceEqual(_bookmarks)), Arg.Any<CancellationToken>());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
[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<TriggerFilter>(), Arg.Is<PageArgs>(x => x.Offset == 0 && x.Limit == 1), Arg.Any<CancellationToken>())
|
||||||
|
.Returns(new Page<StoredTrigger>(firstTriggerPage, 2));
|
||||||
|
_triggerStore.FindManyAsync(Arg.Any<TriggerFilter>(), Arg.Is<PageArgs>(x => x.Offset == 1 && x.Limit == 1), Arg.Any<CancellationToken>())
|
||||||
|
.Returns(new Page<StoredTrigger>(secondTriggerPage, 2));
|
||||||
|
_bookmarkStore.FindManyAsync(Arg.Any<BookmarkFilter>(), Arg.Is<PageArgs>(x => x.Offset == 0 && x.Limit == 1), Arg.Any<CancellationToken>())
|
||||||
|
.Returns(new Page<StoredBookmark>(firstBookmarkPage, 2));
|
||||||
|
_bookmarkStore.FindManyAsync(Arg.Any<BookmarkFilter>(), Arg.Is<PageArgs>(x => x.Offset == 1 && x.Limit == 1), Arg.Any<CancellationToken>())
|
||||||
|
.Returns(new Page<StoredBookmark>(secondBookmarkPage, 2));
|
||||||
|
var task = new CreateSchedulesStartupTask(CreateServiceProvider(), OptionsFactory.Create(_options));
|
||||||
|
|
||||||
|
await task.ExecuteAsync(CancellationToken.None);
|
||||||
|
|
||||||
|
await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is<IEnumerable<StoredTrigger>>(x => x.SequenceEqual(firstTriggerPage)), Arg.Any<CancellationToken>());
|
||||||
|
await _triggerScheduler.Received(1).ScheduleAsync(Arg.Is<IEnumerable<StoredTrigger>>(x => x.SequenceEqual(secondTriggerPage)), Arg.Any<CancellationToken>());
|
||||||
|
await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is<IEnumerable<StoredBookmark>>(x => x.SequenceEqual(firstBookmarkPage)), Arg.Any<CancellationToken>());
|
||||||
|
await _bookmarkScheduler.Received(1).ScheduleAsync(Arg.Is<IEnumerable<StoredBookmark>>(x => x.SequenceEqual(secondBookmarkPage)), Arg.Any<CancellationToken>());
|
||||||
|
}
|
||||||
|
|
||||||
private ServiceProvider CreateServiceProvider(Action<IServiceCollection>? configureServices = null)
|
private ServiceProvider CreateServiceProvider(Action<IServiceCollection>? configureServices = null)
|
||||||
{
|
{
|
||||||
var services = new ServiceCollection();
|
var services = new ServiceCollection();
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue