From 312fcf233a98c673e8d90dd3267acedaa8b80db1 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 22 Nov 2020 18:52:05 +0100 Subject: [PATCH] Fix broken Instant activity --- .../Activities/InstantEvent/InstantEvent.cs | 20 ++++++++- .../HostedServices/TimersHostedService.cs | 2 +- .../Services/IInstantEventManager.cs | 10 ----- .../Services/InstantEventManager.cs | 42 ------------------- .../Triggers/InstantEventTrigger.cs | 8 ++-- .../Triggers/TriggerProviderContext.cs | 2 +- 6 files changed, 25 insertions(+), 59 deletions(-) delete mode 100644 src/activities/Elsa.Activities.Timers/Services/IInstantEventManager.cs delete mode 100644 src/activities/Elsa.Activities.Timers/Services/InstantEventManager.cs diff --git a/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs index 435c9bb23..f4c509aa5 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/InstantEvent/InstantEvent.cs @@ -16,10 +16,28 @@ namespace Elsa.Activities.Timers )] public class InstantEvent : Activity { + private readonly IClock _clock; + + public InstantEvent(IClock clock) + { + _clock = clock; + } + [ActivityProperty(Hint = "An instant in the future at which this activity should execute.")] public Instant Instant { get; set; } + + public Instant ExecuteAt + { + get => GetState(); + set => SetState(value); + } + + protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) + { + ExecuteAt = Instant; + return Suspend(); + } - protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => context.WorkflowExecutionContext.IsFirstPass ? (IActivityExecutionResult)Done() : Suspend(); protected override IActivityExecutionResult OnResume() => Done(); } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs b/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs index e55fbc7ac..7f8586359 100644 --- a/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs +++ b/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs @@ -48,7 +48,7 @@ namespace Elsa.Activities.Timers.HostedServices var now = _clock.GetCurrentInstant(); await workflowScheduler.TriggerWorkflowsAsync(x => x.ExecuteAt <= now, cancellationToken: stoppingToken); await workflowScheduler.TriggerWorkflowsAsync(x => x.ExecuteAt <= now, cancellationToken: stoppingToken); - await workflowScheduler.TriggerWorkflowsAsync(x => x.Instant <= now, cancellationToken: stoppingToken); + await workflowScheduler.TriggerWorkflowsAsync(x => x.ExecuteAt <= now, cancellationToken: stoppingToken); } catch (Exception ex) { diff --git a/src/activities/Elsa.Activities.Timers/Services/IInstantEventManager.cs b/src/activities/Elsa.Activities.Timers/Services/IInstantEventManager.cs deleted file mode 100644 index 2bd910fd7..000000000 --- a/src/activities/Elsa.Activities.Timers/Services/IInstantEventManager.cs +++ /dev/null @@ -1,10 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; - -namespace Elsa.Activities.Timers.Services -{ - public interface IInstantEventManager - { - Task TriggerInstantEventsAsync(CancellationToken cancellationToken); - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/Services/InstantEventManager.cs b/src/activities/Elsa.Activities.Timers/Services/InstantEventManager.cs deleted file mode 100644 index 15801fba0..000000000 --- a/src/activities/Elsa.Activities.Timers/Services/InstantEventManager.cs +++ /dev/null @@ -1,42 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.Activities.Timers.Triggers; -using Elsa.Indexes; -using Elsa.Services; -using Elsa.Triggers; -using NodaTime; -using Open.Linq.AsyncExtensions; - -namespace Elsa.Activities.Timers.Services -{ - public class InstantEventManager : IInstantEventManager - { - private readonly IWorkflowSelector _workflowSelector; - private readonly IWorkflowInstanceManager _workflowInstanceManager; - private readonly IClock _clock; - - public InstantEventManager(IWorkflowSelector workflowSelector, IWorkflowInstanceManager workflowInstanceManager, IClock clock) - { - _workflowSelector = workflowSelector; - _workflowInstanceManager = workflowInstanceManager; - _clock = clock; - } - - public async Task TriggerInstantEventsAsync(CancellationToken cancellationToken) - { - var now = _clock.GetCurrentInstant(); - var candidates = await _workflowSelector.SelectWorkflowsAsync(x => x.Instant <= now, cancellationToken).ToList(); - - foreach (var candidate in candidates) - { - // Only trigger workflows that haven't already executed. - var instanceCount = await _workflowInstanceManager.Query(x => x.WorkflowDefinitionId == candidate.WorkflowBlueprint.Id).CountAsync(); - - if(instanceCount > 0) - continue; - - - } - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs b/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs index a52a93ce7..2c3d31d8f 100644 --- a/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs +++ b/src/activities/Elsa.Activities.Timers/Triggers/InstantEventTrigger.cs @@ -13,7 +13,7 @@ namespace Elsa.Activities.Timers.Triggers /// /// The instant at which the event should trigger. /// - public Instant Instant { get; set; } + public Instant ExecuteAt { get; set; } public override bool IsOneOff => true; } @@ -34,16 +34,16 @@ namespace Elsa.Activities.Timers.Triggers // Only provide a trigger if the workflow hasn't executed already sometime in the past. var workflowDefinitionId = context.ActivityExecutionContext.WorkflowExecutionContext.WorkflowBlueprint.Id; var instanceCount = await _workflowInstanceManager.ListByDefinitionAsync(workflowDefinitionId, cancellationToken).Count(); - var configuredInstant = await context.Activity.GetPropertyValueAsync(x => x.Instant, cancellationToken); + var executeAt = context.Activity.GetState(x => x.ExecuteAt); var now = _clock.GetCurrentInstant(); // If the configured instant lies in the past, and the workflow was already executed once, we don't trigger again. - if (configuredInstant <= now && instanceCount > 0) + if (executeAt <= now && instanceCount > 0) return NullTrigger.Instance; return new InstantEventTrigger { - Instant = await context.Activity.GetPropertyValueAsync(x => x.Instant, cancellationToken) + ExecuteAt = executeAt }; } } diff --git a/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs b/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs index b5332993f..927866bfe 100644 --- a/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs +++ b/src/core/Elsa.Abstractions/Triggers/TriggerProviderContext.cs @@ -11,7 +11,7 @@ namespace Elsa.Triggers } public ActivityExecutionContext ActivityExecutionContext { get; } - public ActivityBlueprintWrapper GetActivity() where TActivity : IActivity => new ActivityBlueprintWrapper(ActivityExecutionContext); + public ActivityBlueprintWrapper GetActivity() where TActivity : IActivity => new(ActivityExecutionContext); } public class TriggerProviderContext : TriggerProviderContext where T:IActivity