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; }