Fix broken Instant activity
This commit is contained in:
parent
239cb54772
commit
312fcf233a
|
|
@ -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<Instant>();
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
|
@ -48,7 +48,7 @@ namespace Elsa.Activities.Timers.HostedServices
|
|||
var now = _clock.GetCurrentInstant();
|
||||
await workflowScheduler.TriggerWorkflowsAsync<TimerEventTrigger>(x => x.ExecuteAt <= now, cancellationToken: stoppingToken);
|
||||
await workflowScheduler.TriggerWorkflowsAsync<CronEventTrigger>(x => x.ExecuteAt <= now, cancellationToken: stoppingToken);
|
||||
await workflowScheduler.TriggerWorkflowsAsync<InstantEventTrigger>(x => x.Instant <= now, cancellationToken: stoppingToken);
|
||||
await workflowScheduler.TriggerWorkflowsAsync<InstantEventTrigger>(x => x.ExecuteAt <= now, cancellationToken: stoppingToken);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -1,10 +0,0 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Elsa.Activities.Timers.Services
|
||||
{
|
||||
public interface IInstantEventManager
|
||||
{
|
||||
Task TriggerInstantEventsAsync(CancellationToken cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<InstantEventTrigger>(x => x.Instant <= now, cancellationToken).ToList();
|
||||
|
||||
foreach (var candidate in candidates)
|
||||
{
|
||||
// Only trigger workflows that haven't already executed.
|
||||
var instanceCount = await _workflowInstanceManager.Query<WorkflowInstanceIndex>(x => x.WorkflowDefinitionId == candidate.WorkflowBlueprint.Id).CountAsync();
|
||||
|
||||
if(instanceCount > 0)
|
||||
continue;
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -13,7 +13,7 @@ namespace Elsa.Activities.Timers.Triggers
|
|||
/// <summary>
|
||||
/// The instant at which the event should trigger.
|
||||
/// </summary>
|
||||
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
|
||||
};
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ namespace Elsa.Triggers
|
|||
}
|
||||
|
||||
public ActivityExecutionContext ActivityExecutionContext { get; }
|
||||
public ActivityBlueprintWrapper<TActivity> GetActivity<TActivity>() where TActivity : IActivity => new ActivityBlueprintWrapper<TActivity>(ActivityExecutionContext);
|
||||
public ActivityBlueprintWrapper<TActivity> GetActivity<TActivity>() where TActivity : IActivity => new(ActivityExecutionContext);
|
||||
}
|
||||
|
||||
public class TriggerProviderContext<T> : TriggerProviderContext where T:IActivity
|
||||
|
|
|
|||
Loading…
Reference in a new issue