diff --git a/src/modules/Elsa.Hangfire/Extensions/JobStorageExtensions.cs b/src/modules/Elsa.Hangfire/Extensions/JobStorageExtensions.cs
index 4dce00947..64ed51057 100644
--- a/src/modules/Elsa.Hangfire/Extensions/JobStorageExtensions.cs
+++ b/src/modules/Elsa.Hangfire/Extensions/JobStorageExtensions.cs
@@ -1,3 +1,4 @@
+using Elsa.Hangfire.Jobs;
using Hangfire;
using Hangfire.Storage.Monitoring;
@@ -11,22 +12,44 @@ public static class JobStorageExtensions
///
/// Enumerates all scheduled jobs of a given type.
///
- public static IEnumerable> EnumerateScheduledJobs(this JobStorage storage, string name)
+ public static IEnumerable> EnumerateScheduledJobs(this JobStorage storage, string name)
{
var api = storage.GetMonitoringApi();
var skip = 0;
const int take = 100;
- JobList jobList;
+ JobList scheduledJobs;
do
{
- jobList = api.ScheduledJobs(skip, take);
+ scheduledJobs = api.ScheduledJobs(skip, take);
- var jobs = jobList.FindAll(x => x.Value.Job.Type == typeof(TJob));
+ var jobs = scheduledJobs.FindAll(x => x.Value.Job.Type == typeof(RunWorkflowJob) || x.Value.Job.Type == typeof(ResumeWorkflowJob));
foreach (var job in jobs.Where(x => (string)x.Value.Job.Args[0] == name))
yield return job;
skip += take;
- } while (jobList.Count == take);
+ } while (scheduledJobs.Count == take);
+ }
+
+ ///
+ /// Enumerates all enqueued jobs of a given type.
+ ///
+ public static IEnumerable> EnumerateQueuedJobs(this JobStorage storage, string queueName, string taskName)
+ {
+ var api = storage.GetMonitoringApi();
+ var skip = 0;
+ const int take = 100;
+ JobList enqueuedJobs;
+
+ do
+ {
+ enqueuedJobs = api.EnqueuedJobs(queueName, skip, take);
+
+ var jobs = enqueuedJobs.FindAll(x => x.Value.Job.Type == typeof(RunWorkflowJob) || x.Value.Job.Type == typeof(ResumeWorkflowJob));
+ foreach (var job in jobs.Where(x => (string)x.Value.Job.Args[0] == taskName))
+ yield return job;
+
+ skip += take;
+ } while (enqueuedJobs.Count == take);
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Services/HangfireWorkflowScheduler.cs b/src/modules/Elsa.Hangfire/Services/HangfireWorkflowScheduler.cs
index 263a586fc..f249895e1 100644
--- a/src/modules/Elsa.Hangfire/Services/HangfireWorkflowScheduler.cs
+++ b/src/modules/Elsa.Hangfire/Services/HangfireWorkflowScheduler.cs
@@ -75,16 +75,27 @@ public class HangfireWorkflowScheduler : IWorkflowScheduler
private void DeleteJobByTaskName(string taskName)
{
- var scheduledJobIds = GetScheduledJobIds(taskName);
+ var scheduledJobIds = GetScheduledJobIds(taskName);
foreach (var jobId in scheduledJobIds) _backgroundJobClient.Delete(jobId);
+ var queuedJobsIds = GetQueuedJobIds(taskName);
+ foreach (var jobId in queuedJobsIds) _backgroundJobClient.Delete(jobId);
+
var recurringJobIds = GetRecurringJobIds(taskName);
foreach (var jobId in recurringJobIds) _recurringJobManager.RemoveIfExists(jobId);
}
- private IEnumerable GetScheduledJobIds(string taskName)
+ private IEnumerable GetScheduledJobIds(string taskName)
{
- return _jobStorage.EnumerateScheduledJobs(taskName)
+ return _jobStorage.EnumerateScheduledJobs(taskName)
+ .Select(x => x.Key)
+ .Distinct()
+ .ToList();
+ }
+
+ private IEnumerable GetQueuedJobIds(string taskName)
+ {
+ return _jobStorage.EnumerateQueuedJobs("default", taskName)
.Select(x => x.Key)
.Distinct()
.ToList();