diff --git a/src/activities/Elsa.Activities.Temporal.Common/StartupTasks/StartJobs.cs b/src/activities/Elsa.Activities.Temporal.Common/StartupTasks/StartJobs.cs index 1484f3e6b..7b5e058fa 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/StartupTasks/StartJobs.cs +++ b/src/activities/Elsa.Activities.Temporal.Common/StartupTasks/StartJobs.cs @@ -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; diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Models/RunHangfireWorkflowJobModel.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Models/RunHangfireWorkflowJobModel.cs index c66b1b310..bad9495b6 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Models/RunHangfireWorkflowJobModel.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Models/RunHangfireWorkflowJobModel.cs @@ -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); + } } } diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs index 27509c53d..93d34bddf 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/HangfireWorkflowScheduler.cs @@ -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; } diff --git a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/JobManager.cs b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/JobManager.cs index 62062010a..f78c3fa3d 100644 --- a/src/activities/Elsa.Activities.Temporal.Hangfire/Services/JobManager.cs +++ b/src/activities/Elsa.Activities.Temporal.Hangfire/Services/JobManager.cs @@ -24,48 +24,51 @@ namespace Elsa.Activities.Temporal.Hangfire.Services public void ScheduleJob(RunHangfireWorkflowJobModel data, Instant instant) { - UnscheduleJob(data); _backgroundJobClient.Schedule(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(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 QueryJobs() => _jobStorage.GetConnection().GetRecurringJobs().Where(x => x.Job.Type == typeof(RunHangfireWorkflowJob)); + private IEnumerable QueryRecurringJobs() => _jobStorage.GetConnection().GetRecurringJobs().Where(x => x.Job.Type == typeof(RunHangfireWorkflowJob)); + } + + public class JobPayload + { + public IDictionary Jobs { get; set; } = new Dictionary(); } } \ No newline at end of file