Fix identity too long for Hangfire hash

This commit is contained in:
Sipke Schoorstra 2021-05-23 20:09:59 +02:00
parent c95f77d9a4
commit 63434ff803
4 changed files with 38 additions and 24 deletions

View file

@ -13,7 +13,7 @@ namespace Elsa.Activities.Temporal.Common.StartupTasks
public class StartJobs : IStartupTask
{
// TODO: Figure out how to start jobs across multiple tenants / how to get a list of all tenants.
private const string TenantId = default;
private const string? TenantId = default;
private readonly IBookmarkFinder _bookmarkFinder;
private readonly IWorkflowScheduler _workflowScheduler;

View file

@ -1,3 +1,7 @@
using System;
using System.Security.Cryptography;
using System.Text;
namespace Elsa.Activities.Temporal.Hangfire.Models
{
public class RunHangfireWorkflowJobModel
@ -18,6 +22,14 @@ namespace Elsa.Activities.Temporal.Hangfire.Models
public string? CronExpression { get; set; }
public bool IsRecurringJob => string.IsNullOrEmpty(CronExpression) == false;
public string GetIdentity() => $"Elsa-tenant:{TenantId ?? "default"}-workflow-instance:{WorkflowInstanceId ?? WorkflowDefinitionId}-activity:{ActivityId}";
public string GetIdentity()
{
var text = $"{TenantId ?? "default"}:{WorkflowInstanceId ?? WorkflowDefinitionId}:{ActivityId}";
var bytes = Encoding.UTF8.GetBytes(text);
using var sha1 = new SHA1Managed();
var hash = sha1.ComputeHash(bytes);
return Convert.ToBase64String(hash);
}
}
}

View file

@ -25,7 +25,7 @@ namespace Elsa.Activities.Temporal.Hangfire.Services
_jobManager.ScheduleJob(data, startAt);
if (cron != null)
_jobManager.ScheduleJob(data, cron);
_jobManager.ScheduleRecurringJob(data, cron);
return Task.CompletedTask;
}
@ -34,7 +34,7 @@ namespace Elsa.Activities.Temporal.Hangfire.Services
{
var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, cronExpression);
_jobManager.ScheduleJob(data, cronExpression);
_jobManager.ScheduleRecurringJob(data, cronExpression);
return Task.CompletedTask;
}
@ -42,14 +42,13 @@ namespace Elsa.Activities.Temporal.Hangfire.Services
public Task UnscheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, CancellationToken cancellationToken = default)
{
var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId);
var identity = data.GetIdentity();
_jobManager.UnscheduleJob(identity);
_jobManager.DeleteJob(data);
return Task.CompletedTask;
}
public Task UnscheduleWorkflowDefinitionAsync(string workflowDefinitionId, string? tenantId, CancellationToken cancellationToken = default)
{
_jobManager.UnscheduleJobs(workflowDefinitionId, tenantId);
_jobManager.DeleteJobs(workflowDefinitionId, tenantId);
return Task.CompletedTask;
}

View file

@ -24,48 +24,51 @@ namespace Elsa.Activities.Temporal.Hangfire.Services
public void ScheduleJob(RunHangfireWorkflowJobModel data, Instant instant)
{
UnscheduleJob(data);
_backgroundJobClient.Schedule<RunHangfireWorkflowJob>(job => job.ExecuteAsync(data), instant.ToDateTimeOffset());
}
public void ScheduleJob(RunHangfireWorkflowJobModel data, string cronExpression)
public void ScheduleRecurringJob(RunHangfireWorkflowJobModel data, string cronExpression)
{
var identity = data.GetIdentity();
_recurringJobManager.AddOrUpdate<RunHangfireWorkflowJob>(identity, job => job.ExecuteAsync(data), cronExpression);
}
public void UnscheduleJob(RunHangfireWorkflowJobModel data)
public void DeleteJob(RunHangfireWorkflowJobModel data)
{
DeleteRecurringJob(data);
}
public void DeleteRecurringJob(RunHangfireWorkflowJobModel data)
{
var identity = data.GetIdentity();
UnscheduleJob(identity);
_recurringJobManager.RemoveIfExists(identity);
}
public void UnscheduleJob(string identity)
public void DeleteJobs(string workflowDefinitionId, string? tenantId)
{
var job = QueryJobs().FirstOrDefault(x => GetJobModel(x.Job).GetIdentity() == identity);
if (job == null)
return;
_backgroundJobClient.Delete(job.Id);
_recurringJobManager.RemoveIfExists(job.Id);
DeleteRecurringJobs(workflowDefinitionId, tenantId);
}
public void UnscheduleJobs(string workflowDefinitionId, string? tenantId)
private void DeleteRecurringJobs(string workflowDefinitionId, string? tenantId)
{
var jobs = QueryJobs();
var recurringJobs = QueryRecurringJobs();
var jobsToRemove = jobs.Where(x =>
var recurringJobsToRemove = recurringJobs.Where(x =>
{
var model = GetJobModel(x.Job);
return model.WorkflowDefinitionId == workflowDefinitionId && model.TenantId == tenantId;
}).ToList();
foreach (var job in jobsToRemove)
foreach (var job in recurringJobsToRemove)
_recurringJobManager.RemoveIfExists(job.Id);
}
private RunHangfireWorkflowJobModel GetJobModel(Job job) => (RunHangfireWorkflowJobModel) job.Args[0];
private IEnumerable<RecurringJobDto> QueryJobs() => _jobStorage.GetConnection().GetRecurringJobs().Where(x => x.Job.Type == typeof(RunHangfireWorkflowJob));
private IEnumerable<RecurringJobDto> QueryRecurringJobs() => _jobStorage.GetConnection().GetRecurringJobs().Where(x => x.Job.Type == typeof(RunHangfireWorkflowJob));
}
public class JobPayload
{
public IDictionary<string, string> Jobs { get; set; } = new Dictionary<string, string>();
}
}