Restore activity state persistence and other bug fixes
This commit is contained in:
parent
ea48910a69
commit
ca97f90120
|
|
@ -6,10 +6,7 @@ using NodaTime;
|
|||
// ReSharper disable once CheckNamespace
|
||||
namespace Elsa.Activities.Timers
|
||||
{
|
||||
[ActivityDefinition(
|
||||
Category = "Timers",
|
||||
Description = "Triggers at a specified interval."
|
||||
)]
|
||||
[ActivityDefinition(Category = "Timers", Description = "Triggers at a specified interval.")]
|
||||
public class TimerEvent : Activity
|
||||
{
|
||||
private readonly IClock _clock;
|
||||
|
|
@ -19,11 +16,15 @@ namespace Elsa.Activities.Timers
|
|||
_clock = clock;
|
||||
}
|
||||
|
||||
[ActivityProperty(Hint = "An expression that evaluates to a Duration value")]
|
||||
[ActivityProperty(Hint = "An expression that evaluates to a Duration value.")]
|
||||
public Duration Timeout { get; set; } = default!;
|
||||
|
||||
private Instant? StartTime { get; set; }
|
||||
|
||||
private Instant? StartTime
|
||||
{
|
||||
get => GetState<Instant?>();
|
||||
set => SetState(value);
|
||||
}
|
||||
|
||||
protected override bool OnCanExecute() => StartTime == null || IsExpired();
|
||||
protected override IActivityExecutionResult OnExecute() => ExecuteInternal();
|
||||
protected override IActivityExecutionResult OnResume() => ExecuteInternal();
|
||||
|
|
|
|||
|
|
@ -9,7 +9,7 @@ namespace Elsa.ActivityResults
|
|||
{
|
||||
var activityDefinition = activityExecutionContext.ActivityBlueprint;
|
||||
var blockingActivity = new BlockingActivity(activityDefinition.Id, activityDefinition.Type);
|
||||
activityExecutionContext.WorkflowExecutionContext.BlockingActivities.Add(blockingActivity);
|
||||
activityExecutionContext.WorkflowExecutionContext.WorkflowInstance.BlockingActivities.Add(blockingActivity);
|
||||
activityExecutionContext.WorkflowExecutionContext.Suspend();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,6 +5,8 @@ namespace Elsa.Comparers
|
|||
{
|
||||
public class BlockingActivityEqualityComparer : IEqualityComparer<BlockingActivity>
|
||||
{
|
||||
public static BlockingActivityEqualityComparer Instance { get; } = new BlockingActivityEqualityComparer();
|
||||
|
||||
public bool Equals(BlockingActivity x, BlockingActivity y) => x.ActivityId.Equals(y.ActivityId);
|
||||
public int GetHashCode(BlockingActivity obj) => obj.ActivityId.GetHashCode();
|
||||
}
|
||||
|
|
|
|||
31
src/core/Elsa.Abstractions/Extensions/JObjectExtensions.cs
Normal file
31
src/core/Elsa.Abstractions/Extensions/JObjectExtensions.cs
Normal file
|
|
@ -0,0 +1,31 @@
|
|||
using System;
|
||||
using Newtonsoft.Json;
|
||||
using Newtonsoft.Json.Linq;
|
||||
using NodaTime;
|
||||
using NodaTime.Serialization.JsonNet;
|
||||
|
||||
namespace Elsa
|
||||
{
|
||||
public static class JObjectExtensions
|
||||
{
|
||||
private static readonly JsonSerializer Serializer = new JsonSerializer().ConfigureForNodaTime(DateTimeZoneProviders.Tzdb);
|
||||
|
||||
public static T GetState<T>(this JObject state, string key, Func<T>? defaultValue = null)
|
||||
{
|
||||
var item = state.GetValue(key, StringComparison.OrdinalIgnoreCase);
|
||||
|
||||
if (item == null || item.Type == JTokenType.Null)
|
||||
return defaultValue != null ? defaultValue() : default!;
|
||||
|
||||
return item.ToObject<T>(Serializer)!;
|
||||
}
|
||||
|
||||
public static T GetState<T>(this JObject state, Type type, string key, Func<T>? defaultValue = null)
|
||||
{
|
||||
var item = state.GetValue(key, StringComparison.OrdinalIgnoreCase);
|
||||
return item != null ? (T)item.ToObject(type, Serializer)! : defaultValue != null ? defaultValue() : default!;
|
||||
}
|
||||
|
||||
public static void SetState(this JObject state, string key, object? value) => state[key] = value != null ? JToken.FromObject(value, Serializer) : null;
|
||||
}
|
||||
}
|
||||
|
|
@ -6,11 +6,12 @@ namespace Elsa.Models
|
|||
{
|
||||
public class WorkflowInstance
|
||||
{
|
||||
private HashSet<BlockingActivity> _blockingActivities = new HashSet<BlockingActivity>(BlockingActivityEqualityComparer.Instance);
|
||||
|
||||
public WorkflowInstance()
|
||||
{
|
||||
Variables = new Variables();
|
||||
Activities = new List<ActivityInstance>();
|
||||
BlockingActivities = new HashSet<BlockingActivity>(new BlockingActivityEqualityComparer());
|
||||
ExecutionLog = new List<ExecutionLogEntry>();
|
||||
ScheduledActivities = new Stack<ScheduledActivity>();
|
||||
}
|
||||
|
|
@ -25,7 +26,13 @@ namespace Elsa.Models
|
|||
public Variables Variables { get; set; }
|
||||
public object? Output { get; set; }
|
||||
public ICollection<ActivityInstance> Activities { get; set; }
|
||||
public HashSet<BlockingActivity> BlockingActivities { get; set; }
|
||||
|
||||
public HashSet<BlockingActivity> BlockingActivities
|
||||
{
|
||||
get => _blockingActivities;
|
||||
set => _blockingActivities = new HashSet<BlockingActivity>(value, BlockingActivityEqualityComparer.Instance);
|
||||
}
|
||||
|
||||
public ICollection<ExecutionLogEntry> ExecutionLog { get; set; }
|
||||
public WorkflowFault? Fault { get; set; }
|
||||
public Stack<ScheduledActivity> ScheduledActivities { get; set; }
|
||||
|
|
|
|||
|
|
@ -1,10 +1,13 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Runtime.CompilerServices;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.ActivityResults;
|
||||
using Elsa.Models;
|
||||
using Elsa.Services.Models;
|
||||
using Microsoft.Extensions.Localization;
|
||||
using Newtonsoft.Json.Linq;
|
||||
|
||||
namespace Elsa.Services
|
||||
{
|
||||
|
|
@ -16,6 +19,7 @@ namespace Elsa.Services
|
|||
public string? DisplayName { get; set; }
|
||||
public string? Description { get; set; }
|
||||
public bool PersistWorkflow { get; set; }
|
||||
public JObject Data { get; set; } = default!;
|
||||
|
||||
public ValueTask<bool> CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) =>
|
||||
OnCanExecuteAsync(context, cancellationToken);
|
||||
|
|
@ -76,5 +80,9 @@ namespace Elsa.Services
|
|||
protected CombinedResult Combine(IEnumerable<IActivityExecutionResult> results) => new CombinedResult(results);
|
||||
protected CombinedResult Combine(params IActivityExecutionResult[] results) => new CombinedResult(results);
|
||||
protected FaultResult Fault(LocalizedString message) => new FaultResult(message);
|
||||
|
||||
protected T GetState<T>(Func<T>? defaultValue = null, [CallerMemberName] string name = null!) => Data.GetState(name, defaultValue);
|
||||
protected T GetState<T>(Type type, Func<T>? defaultValue = null, [CallerMemberName] string name = null!) => Data.GetState(type, name, defaultValue);
|
||||
protected void SetState(object? value, [CallerMemberName] string name = null!) => Data.SetState(name, value);
|
||||
}
|
||||
}
|
||||
|
|
@ -2,6 +2,7 @@ using System.Threading;
|
|||
using System.Threading.Tasks;
|
||||
using Elsa.ActivityResults;
|
||||
using Elsa.Services.Models;
|
||||
using Newtonsoft.Json.Linq;
|
||||
|
||||
namespace Elsa.Services
|
||||
{
|
||||
|
|
@ -36,6 +37,11 @@ namespace Elsa.Services
|
|||
/// A value indicating whether the workflow instance will be persisted automatically upon executing this activity.
|
||||
/// </summary>
|
||||
bool PersistWorkflow { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// A data store for the activity to store information that needs to be persisted as part of the workflow instance.
|
||||
/// </summary>
|
||||
JObject Data { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Returns a value of whether the specified activity can execute.
|
||||
|
|
|
|||
|
|
@ -1,7 +1,9 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Models;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
namespace Elsa.Services.Models
|
||||
|
|
@ -27,6 +29,7 @@ namespace Elsa.Services.Models
|
|||
public object? Input { get; }
|
||||
public object? Output { get; set; }
|
||||
public IReadOnlyCollection<string> Outcomes { get; set; }
|
||||
public ActivityInstance ActivityInstance => WorkflowExecutionContext.WorkflowInstance.Activities.First(x => x.Id == ActivityBlueprint.Id);
|
||||
|
||||
public void SetVariable(string name, object? value) => WorkflowExecutionContext.SetVariable(name, value);
|
||||
public object? GetVariable(string name) => WorkflowExecutionContext.GetVariable(name);
|
||||
|
|
@ -44,7 +47,9 @@ namespace Elsa.Services.Models
|
|||
public IActivity ActivateActivity(string activityType, Action<IActivity>? setupActivity = default)
|
||||
{
|
||||
var activityActivator = ServiceProvider.GetRequiredService<IActivityActivator>();
|
||||
return activityActivator.ActivateActivity(activityType, setupActivity);
|
||||
var activity = activityActivator.ActivateActivity(activityType, setupActivity);
|
||||
activity.Data = ActivityInstance.Data;
|
||||
return activity;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -14,7 +14,7 @@ namespace Elsa.Services.Models
|
|||
}
|
||||
|
||||
public WorkflowBlueprint(
|
||||
string definitionId,
|
||||
string id,
|
||||
int version,
|
||||
bool isSingleton,
|
||||
bool isEnabled,
|
||||
|
|
@ -28,7 +28,7 @@ namespace Elsa.Services.Models
|
|||
IEnumerable<IConnection> connections,
|
||||
IActivityPropertyProviders activityPropertyValueProviders)
|
||||
{
|
||||
DefinitionId = definitionId;
|
||||
Id = id;
|
||||
Version = version;
|
||||
IsSingleton = isSingleton;
|
||||
IsEnabled = isEnabled;
|
||||
|
|
@ -44,7 +44,6 @@ namespace Elsa.Services.Models
|
|||
}
|
||||
|
||||
public string Id { get; set; } = default!;
|
||||
public string DefinitionId { get; set; } = default!;
|
||||
public int Version { get; set; }
|
||||
public bool IsSingleton { get; set; }
|
||||
public bool IsEnabled { get; set; }
|
||||
|
|
|
|||
|
|
@ -26,10 +26,6 @@ namespace Elsa.Services.Models
|
|||
public IWorkflowBlueprint WorkflowBlueprint { get; }
|
||||
public IServiceProvider ServiceProvider { get; }
|
||||
public WorkflowInstance WorkflowInstance { get; }
|
||||
|
||||
public HashSet<BlockingActivity> BlockingActivities { get; } =
|
||||
new HashSet<BlockingActivity>(new BlockingActivityEqualityComparer());
|
||||
|
||||
public bool HasScheduledActivities => WorkflowInstance.ScheduledActivities.Any();
|
||||
public IWorkflowFault? WorkflowFault { get; private set; }
|
||||
public bool IsFirstPass { get; private set; }
|
||||
|
|
|
|||
|
|
@ -18,7 +18,7 @@ namespace Elsa.Activities.ControlFlow
|
|||
|
||||
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context)
|
||||
{
|
||||
context.WorkflowExecutionContext.BlockingActivities.Clear();
|
||||
context.WorkflowExecutionContext.WorkflowInstance.BlockingActivities.Clear();
|
||||
context.WorkflowExecutionContext.Complete();
|
||||
|
||||
return Done(OutputValue);
|
||||
|
|
|
|||
|
|
@ -64,10 +64,10 @@ namespace Elsa.Activities.ControlFlow
|
|||
// Remove any inbound blocking activities.
|
||||
var ancestorActivityIds = workflowExecutionContext.GetInboundActivityPath(Id).ToList();
|
||||
var blockingActivities =
|
||||
workflowExecutionContext.BlockingActivities.Where(x => ancestorActivityIds.Contains(x.ActivityId)).ToList();
|
||||
workflowExecutionContext.WorkflowInstance.BlockingActivities.Where(x => ancestorActivityIds.Contains(x.ActivityId)).ToList();
|
||||
|
||||
foreach (var blockingActivity in blockingActivities)
|
||||
workflowExecutionContext.BlockingActivities.Remove(blockingActivity);
|
||||
workflowExecutionContext.WorkflowInstance.BlockingActivities.Remove(blockingActivity);
|
||||
}
|
||||
|
||||
if (!done)
|
||||
|
|
|
|||
|
|
@ -56,7 +56,6 @@ namespace Elsa.Builders
|
|||
public Func<ActivityExecutionContext, CancellationToken, ValueTask<IActivity>> BuildActivityAsync() =>
|
||||
async (context, cancellationToken) =>
|
||||
{
|
||||
//var activity = context.ActivateActivity(context.ActivityBlueprint.Type, SetupActivity);
|
||||
var activity = context.ActivateActivity(context.ActivityBlueprint.Type);
|
||||
activity.Id = ActivityId;
|
||||
activity.Name = Name;
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ namespace Elsa.Services
|
|||
var activityBlueprints = workflowDefinition.Activities.Select(CreateBlueprint).ToDictionary(x => x.Id);
|
||||
|
||||
return new WorkflowBlueprint(
|
||||
workflowDefinition.WorkflowDefinitionVersionId,
|
||||
workflowDefinition.WorkflowDefinitionId,
|
||||
workflowDefinition.Version,
|
||||
workflowDefinition.IsSingleton,
|
||||
workflowDefinition.IsEnabled,
|
||||
|
|
|
|||
|
|
@ -84,8 +84,7 @@ namespace Elsa.Services
|
|||
object? input = default,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
using var scope = _serviceProvider.CreateScope();
|
||||
var workflowExecutionContext = CreateWorkflowExecutionContext(workflowBlueprint, workflowInstance, scope);
|
||||
var workflowExecutionContext = CreateWorkflowExecutionContext(workflowBlueprint, workflowInstance, _serviceProvider);
|
||||
var activity = activityId != null ? workflowBlueprint.GetActivity(activityId) : default;
|
||||
|
||||
switch (workflowExecutionContext.Status)
|
||||
|
|
@ -165,7 +164,7 @@ namespace Elsa.Services
|
|||
if (!await CanExecuteAsync(workflowExecutionContext, activityBlueprint, input, cancellationToken))
|
||||
return;
|
||||
|
||||
workflowExecutionContext.BlockingActivities.RemoveWhere(x => x.ActivityId == activityBlueprint.Id);
|
||||
workflowExecutionContext.WorkflowInstance.BlockingActivities.RemoveWhere(x => x.ActivityId == activityBlueprint.Id);
|
||||
workflowExecutionContext.Resume();
|
||||
workflowExecutionContext.ScheduleActivity(activityBlueprint.Id, input);
|
||||
await RunAsync(workflowExecutionContext, Resume, cancellationToken);
|
||||
|
|
@ -177,11 +176,9 @@ namespace Elsa.Services
|
|||
object? input,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
using var scope = _serviceProvider.CreateScope();
|
||||
|
||||
var activityExecutionContext = new ActivityExecutionContext(
|
||||
workflowExecutionContext,
|
||||
scope.ServiceProvider,
|
||||
_serviceProvider,
|
||||
activityBlueprint,
|
||||
input);
|
||||
|
||||
|
|
@ -199,22 +196,13 @@ namespace Elsa.Services
|
|||
var scheduledActivity = workflowExecutionContext.PopScheduledActivity();
|
||||
var currentActivityId = scheduledActivity.ActivityId;
|
||||
var activityBlueprint = workflowExecutionContext.WorkflowBlueprint.GetActivity(currentActivityId)!;
|
||||
|
||||
using (var scope = _serviceProvider.CreateScope())
|
||||
{
|
||||
var activityExecutionContext = new ActivityExecutionContext(
|
||||
workflowExecutionContext,
|
||||
scope.ServiceProvider,
|
||||
activityBlueprint,
|
||||
scheduledActivity.Input);
|
||||
|
||||
var activity = await activityBlueprint.CreateActivityAsync(activityExecutionContext, cancellationToken);
|
||||
var result = await activityOperation(activityExecutionContext, activity, cancellationToken);
|
||||
await _mediator.Publish(new ActivityExecuting(activityExecutionContext), cancellationToken);
|
||||
await result.ExecuteAsync(activityExecutionContext, cancellationToken);
|
||||
await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken);
|
||||
}
|
||||
|
||||
var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, _serviceProvider, activityBlueprint, scheduledActivity.Input);
|
||||
var activity = await activityBlueprint.CreateActivityAsync(activityExecutionContext, cancellationToken);
|
||||
var result = await activityOperation(activityExecutionContext, activity, cancellationToken);
|
||||
await _mediator.Publish(new ActivityExecuting(activityExecutionContext), cancellationToken);
|
||||
await result.ExecuteAsync(activityExecutionContext, cancellationToken);
|
||||
await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken);
|
||||
|
||||
activityOperation = Execute;
|
||||
workflowExecutionContext.CompletePass();
|
||||
}
|
||||
|
|
@ -223,13 +211,7 @@ namespace Elsa.Services
|
|||
workflowExecutionContext.Complete();
|
||||
}
|
||||
|
||||
private static WorkflowExecutionContext CreateWorkflowExecutionContext(
|
||||
IWorkflowBlueprint workflowBlueprint,
|
||||
WorkflowInstance workflowInstance,
|
||||
IServiceScope serviceScope) =>
|
||||
new WorkflowExecutionContext(
|
||||
serviceScope.ServiceProvider,
|
||||
workflowBlueprint,
|
||||
workflowInstance);
|
||||
private static WorkflowExecutionContext CreateWorkflowExecutionContext(IWorkflowBlueprint workflowBlueprint, WorkflowInstance workflowInstance, IServiceProvider serviceProvider) =>
|
||||
new WorkflowExecutionContext(serviceProvider, workflowBlueprint, workflowInstance);
|
||||
}
|
||||
}
|
||||
|
|
@ -151,9 +151,7 @@ namespace Elsa.Services
|
|||
else
|
||||
predicate = x => x.ActivityType == activityType;
|
||||
|
||||
var query = _workflowInstanceManager.Query<WorkflowInstanceBlockingActivitiesIndex>()
|
||||
.Where(predicate);
|
||||
|
||||
var query = _workflowInstanceManager.Query<WorkflowInstanceBlockingActivitiesIndex>().Where(predicate);
|
||||
var workflowInstances = await query.ListAsync();
|
||||
var tuples = workflowInstances.GetBlockingActivities();
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue