Fix Hangfire provider job cancellation cleanup

This commit is contained in:
Sipke Schoorstra 2023-04-13 22:34:07 +02:00
parent fb02379133
commit 1a389afb81
2 changed files with 42 additions and 8 deletions

View file

@ -1,3 +1,4 @@
using Elsa.Hangfire.Jobs;
using Hangfire;
using Hangfire.Storage.Monitoring;
@ -11,22 +12,44 @@ public static class JobStorageExtensions
/// <summary>
/// Enumerates all scheduled jobs of a given type.
/// </summary>
public static IEnumerable<KeyValuePair<string, ScheduledJobDto>> EnumerateScheduledJobs<TJob>(this JobStorage storage, string name)
public static IEnumerable<KeyValuePair<string, ScheduledJobDto>> EnumerateScheduledJobs(this JobStorage storage, string name)
{
var api = storage.GetMonitoringApi();
var skip = 0;
const int take = 100;
JobList<ScheduledJobDto> jobList;
JobList<ScheduledJobDto> 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);
}
/// <summary>
/// Enumerates all enqueued jobs of a given type.
/// </summary>
public static IEnumerable<KeyValuePair<string, EnqueuedJobDto>> EnumerateQueuedJobs(this JobStorage storage, string queueName, string taskName)
{
var api = storage.GetMonitoringApi();
var skip = 0;
const int take = 100;
JobList<EnqueuedJobDto> 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);
}
}

View file

@ -75,16 +75,27 @@ public class HangfireWorkflowScheduler : IWorkflowScheduler
private void DeleteJobByTaskName(string taskName)
{
var scheduledJobIds = GetScheduledJobIds<RunWorkflowJob>(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<RunWorkflowJob>(taskName);
foreach (var jobId in recurringJobIds) _recurringJobManager.RemoveIfExists(jobId);
}
private IEnumerable<string> GetScheduledJobIds<TJob>(string taskName)
private IEnumerable<string> GetScheduledJobIds(string taskName)
{
return _jobStorage.EnumerateScheduledJobs<TJob>(taskName)
return _jobStorage.EnumerateScheduledJobs(taskName)
.Select(x => x.Key)
.Distinct()
.ToList();
}
private IEnumerable<string> GetQueuedJobIds(string taskName)
{
return _jobStorage.EnumerateQueuedJobs("default", taskName)
.Select(x => x.Key)
.Distinct()
.ToList();