Incremental work on timer event
This commit is contained in:
parent
26ee6b10a7
commit
dce2c4f9db
|
|
@ -24,10 +24,10 @@ namespace Elsa.Activities.Timers
|
|||
this.clock = clock;
|
||||
}
|
||||
|
||||
[ActivityProperty(Hint = "An expression that evaluates to a TimeSpan value")]
|
||||
public IWorkflowExpression<TimeSpan> Timeout
|
||||
[ActivityProperty(Hint = "An expression that evaluates to a Duration value")]
|
||||
public IWorkflowExpression<Duration> Timeout
|
||||
{
|
||||
get => GetState<IWorkflowExpression<TimeSpan>>(() => new LiteralExpression<TimeSpan>("00:01:00"));
|
||||
get => GetState<IWorkflowExpression<Duration>>(() => new LiteralExpression<Duration>("00:01:00"));
|
||||
set => SetState(value);
|
||||
}
|
||||
|
||||
|
|
@ -42,9 +42,10 @@ namespace Elsa.Activities.Timers
|
|||
return StartTime == null || await IsExpiredAsync(context, cancellationToken);
|
||||
}
|
||||
|
||||
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => Suspend();
|
||||
protected override Task<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => ExecuteInternalAsync(context, cancellationToken);
|
||||
protected override Task<IActivityExecutionResult> OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => ExecuteInternalAsync(context, cancellationToken);
|
||||
|
||||
protected override async Task<IActivityExecutionResult> OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
|
||||
private async Task<IActivityExecutionResult> ExecuteInternalAsync(ActivityExecutionContext context, CancellationToken cancellationToken)
|
||||
{
|
||||
if (await IsExpiredAsync(context, cancellationToken))
|
||||
{
|
||||
|
|
@ -61,11 +62,12 @@ namespace Elsa.Activities.Timers
|
|||
|
||||
if (StartTime == null)
|
||||
StartTime = now;
|
||||
|
||||
var startTime = StartTime.Value;
|
||||
var timeout = await context.EvaluateAsync(Timeout, cancellationToken);
|
||||
var expiresAt = startTime + timeout - Duration.FromMilliseconds(1000);
|
||||
|
||||
var timeSpan = await context.EvaluateAsync(Timeout, cancellationToken);
|
||||
var expiresAt = StartTime.Value.ToDateTimeUtc() + timeSpan;
|
||||
|
||||
return now.ToDateTimeUtc() >= expiresAt;
|
||||
return now >= expiresAt;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -11,9 +11,9 @@ namespace Elsa.Activities.MassTransit
|
|||
public static class TimerEventBuilderExtensions
|
||||
{
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, Action<TimerEvent>? setup = default) => builder.Then(setup);
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, IWorkflowExpression<TimeSpan> value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, Func<ActivityExecutionContext, TimeSpan> value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, Func<TimeSpan> value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, TimeSpan value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, IWorkflowExpression<Duration> value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, Func<ActivityExecutionContext, Duration> value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, Func<Duration> value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
public static ActivityBuilder TimerEvent(this IBuilder builder, Duration value) => builder.TimerEvent(x => x.WithTimeout(value));
|
||||
}
|
||||
}
|
||||
|
|
@ -10,9 +10,9 @@ namespace Elsa.Activities.MassTransit
|
|||
{
|
||||
public static class TimerEventExtensions
|
||||
{
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, IWorkflowExpression<TimeSpan> value) => activity.With(x => x.Timeout, value);
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, Func<ActivityExecutionContext, TimeSpan> value) => activity.With(x => x.Timeout, new CodeExpression<TimeSpan>(value));
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, Func<TimeSpan> value) => activity.With(x => x.Timeout, new CodeExpression<TimeSpan>(value));
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, TimeSpan value) => activity.With(x => x.Timeout, new CodeExpression<TimeSpan>(value));
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, IWorkflowExpression<Duration> value) => activity.With(x => x.Timeout, value);
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, Func<ActivityExecutionContext, Duration> value) => activity.With(x => x.Timeout, new CodeExpression<Duration>(value));
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, Func<Duration> value) => activity.With(x => x.Timeout, new CodeExpression<Duration>(value));
|
||||
public static TimerEvent WithTimeout(this TimerEvent activity, Duration value) => activity.With(x => x.Timeout, new CodeExpression<Duration>(value));
|
||||
}
|
||||
}
|
||||
|
|
@ -63,7 +63,7 @@ namespace Elsa.Services
|
|||
logger.LogDebug("Workflow instance {WorkflowInstanceId} does not exist.", workflowInstanceId);
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
var workflow = await workflowRegistry.GetWorkflowAsync(workflowInstance.DefinitionId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken);
|
||||
return await RunAsync(workflow, workflowInstance, activityId, input, cancellationToken);
|
||||
}
|
||||
|
|
@ -72,10 +72,10 @@ namespace Elsa.Services
|
|||
{
|
||||
var workflow = await workflowRegistry.GetWorkflowAsync(workflowDefinitionId, VersionOptions.Published, cancellationToken);
|
||||
var workflowInstance = await workflowActivator.ActivateAsync(workflow, correlationId, cancellationToken);
|
||||
|
||||
|
||||
return await RunAsync(workflow, workflowInstance, activityId, input, cancellationToken);
|
||||
}
|
||||
|
||||
|
||||
public async Task<WorkflowExecutionContext> RunWorkflowAsync(Workflow workflow, string? activityId = default, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var workflowInstance = await workflowActivator.ActivateAsync(workflow, correlationId, cancellationToken);
|
||||
|
|
@ -100,6 +100,29 @@ namespace Elsa.Services
|
|||
break;
|
||||
}
|
||||
|
||||
await mediator.Publish(new WorkflowExecuted(workflowExecutionContext), cancellationToken);
|
||||
|
||||
var statusEvent = default(object);
|
||||
|
||||
switch (workflowExecutionContext.Status)
|
||||
{
|
||||
case WorkflowStatus.Cancelled:
|
||||
statusEvent = new WorkflowCancelled(workflowExecutionContext);
|
||||
break;
|
||||
case WorkflowStatus.Completed:
|
||||
statusEvent = new WorkflowCompleted(workflowExecutionContext);
|
||||
break;
|
||||
case WorkflowStatus.Faulted:
|
||||
statusEvent = new WorkflowFaulted(workflowExecutionContext);
|
||||
break;
|
||||
case WorkflowStatus.Suspended:
|
||||
statusEvent = new WorkflowSuspended(workflowExecutionContext);
|
||||
break;
|
||||
}
|
||||
|
||||
if (statusEvent != null)
|
||||
await mediator.Publish(statusEvent, cancellationToken);
|
||||
|
||||
return workflowExecutionContext;
|
||||
}
|
||||
|
||||
|
|
@ -107,12 +130,15 @@ namespace Elsa.Services
|
|||
{
|
||||
if (activity == null)
|
||||
activity = workflowExecutionContext.GetStartActivities().First();
|
||||
|
||||
|
||||
if (!await CanExecuteAsync(workflowExecutionContext, activity, input, cancellationToken))
|
||||
return;
|
||||
|
||||
workflowExecutionContext.Status = WorkflowStatus.Running;
|
||||
workflowExecutionContext.ScheduleActivity(activity, input);
|
||||
await RunAsync(workflowExecutionContext, Execute, cancellationToken);
|
||||
}
|
||||
|
||||
|
||||
private async Task RunWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken)
|
||||
{
|
||||
await RunAsync(workflowExecutionContext, Execute, cancellationToken);
|
||||
|
|
@ -120,11 +146,20 @@ namespace Elsa.Services
|
|||
|
||||
private async Task ResumeWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, IActivity activity, object? input, CancellationToken cancellationToken)
|
||||
{
|
||||
if (!await CanExecuteAsync(workflowExecutionContext, activity, input, cancellationToken))
|
||||
return;
|
||||
|
||||
workflowExecutionContext.BlockingActivities.Remove(activity);
|
||||
workflowExecutionContext.Status = WorkflowStatus.Running;
|
||||
workflowExecutionContext.ScheduleActivity(activity, input);
|
||||
await RunAsync(workflowExecutionContext, Resume, cancellationToken);
|
||||
}
|
||||
|
||||
private Task<bool> CanExecuteAsync(WorkflowExecutionContext workflowExecutionContext, IActivity activity, object? input, CancellationToken cancellationToken)
|
||||
{
|
||||
var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, activity, Variable.From(input));
|
||||
return activity.CanExecuteAsync(activityExecutionContext, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task RunAsync(
|
||||
WorkflowExecutionContext workflowExecutionContext,
|
||||
|
|
@ -146,29 +181,6 @@ namespace Elsa.Services
|
|||
|
||||
if (workflowExecutionContext.Status == WorkflowStatus.Running)
|
||||
workflowExecutionContext.Complete();
|
||||
|
||||
await mediator.Publish(new WorkflowExecuted(workflowExecutionContext), cancellationToken);
|
||||
|
||||
var statusEvent = default(object);
|
||||
|
||||
switch (workflowExecutionContext.Status)
|
||||
{
|
||||
case WorkflowStatus.Cancelled:
|
||||
statusEvent = new WorkflowCancelled(workflowExecutionContext);
|
||||
break;
|
||||
case WorkflowStatus.Completed:
|
||||
statusEvent = new WorkflowCompleted(workflowExecutionContext);
|
||||
break;
|
||||
case WorkflowStatus.Faulted:
|
||||
statusEvent = new WorkflowFaulted(workflowExecutionContext);
|
||||
break;
|
||||
case WorkflowStatus.Suspended:
|
||||
statusEvent = new WorkflowSuspended(workflowExecutionContext);
|
||||
break;
|
||||
}
|
||||
|
||||
if (statusEvent != null)
|
||||
await mediator.Publish(statusEvent, cancellationToken);
|
||||
}
|
||||
|
||||
private ScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel, IDictionary<string, IActivity> activityLookup)
|
||||
|
|
@ -196,7 +208,7 @@ namespace Elsa.Services
|
|||
activity.State = activityInstance.State;
|
||||
activity.Output = activityInstance.Output;
|
||||
}
|
||||
|
||||
|
||||
return CreateWorkflowExecutionContext(
|
||||
workflowInstance.Id,
|
||||
workflow.DefinitionId,
|
||||
|
|
|
|||
|
|
@ -93,6 +93,7 @@ namespace Elsa.Services
|
|||
var tuples = (IList<(Workflow, IActivity)>)query.ToList();
|
||||
|
||||
tuples = (await FilterRunningSingletonsAsync(tuples, cancellationToken)).ToList();
|
||||
tuples = (await FilterStartedWorkflowsAsync(tuples, cancellationToken)).ToList();
|
||||
|
||||
foreach (var (workflow, activity) in tuples)
|
||||
{
|
||||
|
|
@ -125,24 +126,45 @@ namespace Elsa.Services
|
|||
}
|
||||
|
||||
private async Task<IEnumerable<(Workflow, IActivity)>> FilterRunningSingletonsAsync(
|
||||
IEnumerable<(Workflow, IActivity)> workflows,
|
||||
IEnumerable<(Workflow Workflow, IActivity Activity)> tuples,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
var definitions = workflows.ToList();
|
||||
var transients = definitions.Where(x => !x.Item1.IsSingleton).ToList();
|
||||
var singletons = definitions.Where(x => x.Item1.IsSingleton).ToList();
|
||||
var tupleList = tuples.ToList();
|
||||
var transients = tupleList.Where(x => !x.Workflow.IsSingleton).ToList();
|
||||
var singletons = tupleList.Where(x => x.Workflow.IsSingleton).ToList();
|
||||
var result = transients.ToList();
|
||||
|
||||
foreach (var definition in singletons)
|
||||
foreach (var tuple in singletons)
|
||||
{
|
||||
var instances = await workflowInstanceStore.ListByStatusAsync(
|
||||
definition.Item1.DefinitionId,
|
||||
tuple.Workflow.DefinitionId,
|
||||
WorkflowStatus.Suspended,
|
||||
cancellationToken
|
||||
);
|
||||
|
||||
if (!instances.Any())
|
||||
result.Add(definition);
|
||||
result.Add(tuple);
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
private async Task<IEnumerable<(Workflow, IActivity)>> FilterStartedWorkflowsAsync(
|
||||
IEnumerable<(Workflow Workflow, IActivity Activity)> tuples,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
var tupleList = tuples.ToList();
|
||||
var result = new List<(Workflow, IActivity)>();
|
||||
|
||||
foreach (var tuple in tupleList)
|
||||
{
|
||||
var suspendedInstances = await workflowInstanceStore.ListByStatusAsync(tuple.Workflow.DefinitionId, WorkflowStatus.Suspended, cancellationToken);
|
||||
var idleInstances = await workflowInstanceStore.ListByStatusAsync(tuple.Workflow.DefinitionId, WorkflowStatus.Idle, cancellationToken);
|
||||
var startActivities = tuple.Workflow.GetStartActivities().Select(x => x.Id).ToList();
|
||||
var hasStartedInstances = suspendedInstances.Any(x => x.BlockingActivities.Any(y => startActivities.Contains(y.ActivityId)));
|
||||
|
||||
if(!hasStartedInstances && !idleInstances.Any())
|
||||
result.Add(tuple);
|
||||
}
|
||||
|
||||
return result;
|
||||
|
|
|
|||
|
|
@ -21,7 +21,7 @@ namespace Elsa.Samples.Timers
|
|||
{
|
||||
services
|
||||
.AddElsa()
|
||||
.AddTimerActivities(options => options.Configure(timer => timer.SweepInterval = Duration.FromSeconds(1)))
|
||||
.AddTimerActivities(options => options.Configure(timer => timer.SweepInterval = Duration.FromSeconds(5)))
|
||||
.AddWorkflow<RecurringTaskWorkflow>();
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ using System;
|
|||
using Elsa.Activities.Console;
|
||||
using Elsa.Activities.MassTransit;
|
||||
using Elsa.Builders;
|
||||
using NodaTime;
|
||||
|
||||
namespace Elsa.Samples.Timers
|
||||
{
|
||||
|
|
@ -10,7 +11,7 @@ namespace Elsa.Samples.Timers
|
|||
public void Build(IWorkflowBuilder builder)
|
||||
{
|
||||
builder
|
||||
.TimerEvent(TimeSpan.FromSeconds(5))
|
||||
.TimerEvent(Duration.FromSeconds(5))
|
||||
.WriteLine(() => $"Timer event at {DateTime.Now}");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue