From 1a389afb81feffc4ae3446432a0efbef5837a174 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 13 Apr 2023 22:34:07 +0200 Subject: [PATCH] Fix Hangfire provider job cancellation cleanup --- .../Extensions/JobStorageExtensions.cs | 33 ++++++++++++++++--- .../Services/HangfireWorkflowScheduler.cs | 17 ++++++++-- 2 files changed, 42 insertions(+), 8 deletions(-) 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();