diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Extensions/TimersOptionsExtensions.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Extensions/TimersOptionsExtensions.cs index 9c2ce6ce4..f4c92c0a1 100644 --- a/src/activities/Elsa.Activities.Temporal.Quartz/Extensions/TimersOptionsExtensions.cs +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Extensions/TimersOptionsExtensions.cs @@ -28,6 +28,7 @@ namespace Elsa Action? configureQuartzHostedService = default) { timersOptions.Services + .AddSingleton() .AddSingleton() .AddSingleton() .AddSingleton() diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzSchedulerProvider.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzSchedulerProvider.cs new file mode 100644 index 000000000..dfc9a1c4e --- /dev/null +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzSchedulerProvider.cs @@ -0,0 +1,36 @@ +using System.Threading; +using System.Threading.Tasks; +using Quartz; + +namespace Elsa.Activities.Temporal.Quartz.Services +{ + public class QuartzSchedulerProvider + { + private readonly ISchedulerFactory _schedulerFactory; + private readonly SemaphoreSlim _semaphore = new(1); + private IScheduler? _scheduler; + + public QuartzSchedulerProvider(ISchedulerFactory schedulerFactory) => _schedulerFactory = schedulerFactory; + + public async Task GetSchedulerAsync(CancellationToken cancellationToken) + { + if (_scheduler != null) + return _scheduler; + + await _semaphore.WaitAsync(cancellationToken); + + try + { + if (_scheduler != null) + return _scheduler; + + _scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + return _scheduler!; + } + finally + { + _semaphore.Release(); + } + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs index 8cb5f6eeb..c8d6fd7e3 100644 --- a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs @@ -12,13 +12,13 @@ namespace Elsa.Activities.Temporal.Quartz.Services public class QuartzWorkflowDefinitionScheduler : IWorkflowDefinitionScheduler { private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowDefinitionJob); - private readonly ISchedulerFactory _schedulerFactory; + private readonly QuartzSchedulerProvider _schedulerProvider; private readonly ILogger _logger; private readonly SemaphoreSlim _semaphore = new(1); - public QuartzWorkflowDefinitionScheduler(ISchedulerFactory schedulerFactory, ILogger logger) + public QuartzWorkflowDefinitionScheduler(QuartzSchedulerProvider schedulerProvider, ILogger logger) { - _schedulerFactory = schedulerFactory; + _schedulerProvider = schedulerProvider; _logger = logger; } @@ -41,7 +41,7 @@ namespace Elsa.Activities.Temporal.Quartz.Services public async Task UnscheduleAsync(string workflowDefinitionId, string activityId, CancellationToken cancellationToken) { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var scheduler = await _schedulerProvider.GetSchedulerAsync(cancellationToken); var trigger = CreateTriggerKey(workflowDefinitionId, activityId); var existingTrigger = await scheduler.GetTrigger(trigger, cancellationToken); @@ -51,7 +51,7 @@ namespace Elsa.Activities.Temporal.Quartz.Services public async Task UnscheduleAsync(string workflowDefinitionId, CancellationToken cancellationToken) { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var scheduler = await _schedulerProvider.GetSchedulerAsync(cancellationToken); var groupName = CreateTriggerGroupKey(workflowDefinitionId); var existingTriggers = await scheduler.GetTriggerKeys(GroupMatcher.GroupEquals(groupName), cancellationToken); @@ -61,18 +61,18 @@ namespace Elsa.Activities.Temporal.Quartz.Services public async Task UnscheduleAllAsync(CancellationToken cancellationToken) { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var scheduler = await _schedulerProvider.GetSchedulerAsync(cancellationToken); var jobKeys = await scheduler.GetJobKeys(GroupMatcher.GroupStartsWith("workflow-definition"), cancellationToken); await scheduler.DeleteJobs(jobKeys, cancellationToken); } private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) { + var scheduler = await _schedulerProvider.GetSchedulerAsync(cancellationToken); await _semaphore.WaitAsync(cancellationToken); try { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); var existingTrigger = await scheduler.GetTrigger(trigger.Key, cancellationToken); // For workflow definitions we only schedule the job if one doesn't exist already because another node may have created it beforehand. diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs index 3bee32045..60d9562cc 100644 --- a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs @@ -12,13 +12,13 @@ namespace Elsa.Activities.Temporal.Quartz.Services public class QuartzWorkflowInstanceScheduler : IWorkflowInstanceScheduler { private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowInstanceJob); - private readonly ISchedulerFactory _schedulerFactory; + private readonly QuartzSchedulerProvider _schedulerProvider; private readonly ILogger _logger; private readonly SemaphoreSlim _semaphore = new(1); - public QuartzWorkflowInstanceScheduler(ISchedulerFactory schedulerFactory, ILogger logger) + public QuartzWorkflowInstanceScheduler(QuartzSchedulerProvider schedulerProvider, ILogger logger) { - _schedulerFactory = schedulerFactory; + _schedulerProvider = schedulerProvider; _logger = logger; } @@ -41,7 +41,7 @@ namespace Elsa.Activities.Temporal.Quartz.Services public async Task UnscheduleAsync(string workflowInstanceId, string activityId, CancellationToken cancellationToken) { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var scheduler = await _schedulerProvider. GetSchedulerAsync(cancellationToken); var trigger = CreateTriggerKey(workflowInstanceId, activityId); var existingTrigger = await scheduler.GetTrigger(trigger, cancellationToken); @@ -51,7 +51,7 @@ namespace Elsa.Activities.Temporal.Quartz.Services public async Task UnscheduleAsync(string workflowInstanceId, CancellationToken cancellationToken) { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var scheduler = await _schedulerProvider. GetSchedulerAsync(cancellationToken); var groupName = CreateTriggerGroupKey(workflowInstanceId); var existingTriggers = await scheduler.GetTriggerKeys(GroupMatcher.GroupEquals(groupName), cancellationToken); @@ -61,18 +61,18 @@ namespace Elsa.Activities.Temporal.Quartz.Services public async Task UnscheduleAllAsync(CancellationToken cancellationToken) { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + var scheduler = await _schedulerProvider. GetSchedulerAsync(cancellationToken); var jobKeys = await scheduler.GetJobKeys(GroupMatcher.GroupStartsWith("workflow-instance"), cancellationToken); await scheduler.DeleteJobs(jobKeys, cancellationToken); } private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) { + var scheduler = await _schedulerProvider. GetSchedulerAsync(cancellationToken); await _semaphore.WaitAsync(cancellationToken); try { - var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); var existingTrigger = await scheduler.GetTrigger(trigger.Key, cancellationToken); if (existingTrigger != null) @@ -102,7 +102,7 @@ namespace Elsa.Activities.Temporal.Quartz.Services var groupName = CreateTriggerGroupKey(workflowInstanceId); return new TriggerKey($"activity:{activityId}", groupName); } - + private string CreateTriggerGroupKey(string workflowInstanceId) => $"workflow-instance:{workflowInstanceId}"; } } \ No newline at end of file