From de6fce52d469d396c5b92c6380404ad1eaff2c97 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 9 Feb 2021 21:31:09 +0100 Subject: [PATCH] Re-enable re-queueing of messages for locked workflows and update timer activities to not execute if scheduled time lies in future --- .../Elsa.Activities.Timers/Activities/Cron/Cron.cs | 7 +++++++ .../Elsa.Activities.Timers/Activities/StartAt/StartAt.cs | 7 +++++++ .../Elsa.Activities.Timers/Activities/Timer/Timer.cs | 7 +++++++ .../Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs | 5 ++--- 4 files changed, 23 insertions(+), 3 deletions(-) diff --git a/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs b/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs index 1b8b9251c..50c4dba37 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/Cron/Cron.cs @@ -38,6 +38,13 @@ namespace Elsa.Activities.Timers get => GetState(); set => SetState(value); } + + protected override bool OnCanExecute(ActivityExecutionContext context) + { + var now = _clock.GetCurrentInstant(); + var executeAt = ExecuteAt; + return executeAt == null || executeAt <= now; + } protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) { diff --git a/src/activities/Elsa.Activities.Timers/Activities/StartAt/StartAt.cs b/src/activities/Elsa.Activities.Timers/Activities/StartAt/StartAt.cs index b32893755..b2138b08c 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/StartAt/StartAt.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/StartAt/StartAt.cs @@ -36,6 +36,13 @@ namespace Elsa.Activities.Timers set => SetState(value); } + protected override bool OnCanExecute(ActivityExecutionContext context) + { + var executeAt = ExecuteAt; + var now = _clock.GetCurrentInstant(); + return executeAt == null || executeAt <= now; + } + protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { if (context.WorkflowExecutionContext.IsFirstPass) diff --git a/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs b/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs index 64b7fa615..0d706caca 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/Timer/Timer.cs @@ -29,6 +29,13 @@ namespace Elsa.Activities.Timers get => GetState(); set => SetState(value); } + + protected override bool OnCanExecute(ActivityExecutionContext context) + { + var now = _clock.GetCurrentInstant(); + var executeAt = ExecuteAt; + return executeAt == null || executeAt <= now; + } protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs index e724e53c1..49dc7c223 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs @@ -46,9 +46,8 @@ namespace Elsa.Consumers if (!await _distributedLockProvider.AcquireLockAsync(lockKey)) { - // TODO: Reschedule message if it's not a redelivery. - // var currentContext = MessageContext.Current; - _logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}", workflowInstanceId); + _logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Re-queueing message", workflowInstanceId); + await _commandSender.SendAsync(message); return; }