diff --git a/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs index 2df6bf797..cbf39ab60 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs @@ -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(); + set => SetState(value); + } + protected override bool OnCanExecute() => StartTime == null || IsExpired(); protected override IActivityExecutionResult OnExecute() => ExecuteInternal(); protected override IActivityExecutionResult OnResume() => ExecuteInternal(); diff --git a/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs b/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs index 54f01cf32..cf6e7233f 100644 --- a/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs +++ b/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs @@ -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(); } } diff --git a/src/core/Elsa.Abstractions/Comparers/BlockingActivityEqualityComparer.cs b/src/core/Elsa.Abstractions/Comparers/BlockingActivityEqualityComparer.cs index e3c773874..27fb9fc7d 100644 --- a/src/core/Elsa.Abstractions/Comparers/BlockingActivityEqualityComparer.cs +++ b/src/core/Elsa.Abstractions/Comparers/BlockingActivityEqualityComparer.cs @@ -5,6 +5,8 @@ namespace Elsa.Comparers { public class BlockingActivityEqualityComparer : IEqualityComparer { + 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(); } diff --git a/src/core/Elsa.Abstractions/Extensions/JObjectExtensions.cs b/src/core/Elsa.Abstractions/Extensions/JObjectExtensions.cs new file mode 100644 index 000000000..592faaa97 --- /dev/null +++ b/src/core/Elsa.Abstractions/Extensions/JObjectExtensions.cs @@ -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(this JObject state, string key, Func? defaultValue = null) + { + var item = state.GetValue(key, StringComparison.OrdinalIgnoreCase); + + if (item == null || item.Type == JTokenType.Null) + return defaultValue != null ? defaultValue() : default!; + + return item.ToObject(Serializer)!; + } + + public static T GetState(this JObject state, Type type, string key, Func? 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; + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs index 5ee8c3a65..17bd506c5 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs @@ -6,11 +6,12 @@ namespace Elsa.Models { public class WorkflowInstance { + private HashSet _blockingActivities = new HashSet(BlockingActivityEqualityComparer.Instance); + public WorkflowInstance() { Variables = new Variables(); Activities = new List(); - BlockingActivities = new HashSet(new BlockingActivityEqualityComparer()); ExecutionLog = new List(); ScheduledActivities = new Stack(); } @@ -25,7 +26,13 @@ namespace Elsa.Models public Variables Variables { get; set; } public object? Output { get; set; } public ICollection Activities { get; set; } - public HashSet BlockingActivities { get; set; } + + public HashSet BlockingActivities + { + get => _blockingActivities; + set => _blockingActivities = new HashSet(value, BlockingActivityEqualityComparer.Instance); + } + public ICollection ExecutionLog { get; set; } public WorkflowFault? Fault { get; set; } public Stack ScheduledActivities { get; set; } diff --git a/src/core/Elsa.Abstractions/Services/Activity.cs b/src/core/Elsa.Abstractions/Services/Activity.cs index 4e826f113..25a606441 100644 --- a/src/core/Elsa.Abstractions/Services/Activity.cs +++ b/src/core/Elsa.Abstractions/Services/Activity.cs @@ -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 CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(context, cancellationToken); @@ -76,5 +80,9 @@ namespace Elsa.Services protected CombinedResult Combine(IEnumerable results) => new CombinedResult(results); protected CombinedResult Combine(params IActivityExecutionResult[] results) => new CombinedResult(results); protected FaultResult Fault(LocalizedString message) => new FaultResult(message); + + protected T GetState(Func? defaultValue = null, [CallerMemberName] string name = null!) => Data.GetState(name, defaultValue); + protected T GetState(Type type, Func? defaultValue = null, [CallerMemberName] string name = null!) => Data.GetState(type, name, defaultValue); + protected void SetState(object? value, [CallerMemberName] string name = null!) => Data.SetState(name, value); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IActivity.cs b/src/core/Elsa.Abstractions/Services/IActivity.cs index df1da1f25..3eb111a84 100644 --- a/src/core/Elsa.Abstractions/Services/IActivity.cs +++ b/src/core/Elsa.Abstractions/Services/IActivity.cs @@ -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. /// bool PersistWorkflow { get; set; } + + /// + /// A data store for the activity to store information that needs to be persisted as part of the workflow instance. + /// + JObject Data { get; set; } /// /// Returns a value of whether the specified activity can execute. diff --git a/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs index 1b6b7f543..0e96fa362 100644 --- a/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs @@ -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 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? setupActivity = default) { var activityActivator = ServiceProvider.GetRequiredService(); - return activityActivator.ActivateActivity(activityType, setupActivity); + var activity = activityActivator.ActivateActivity(activityType, setupActivity); + activity.Data = ActivityInstance.Data; + return activity; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs index 81d3e3f9d..cd4792ac6 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs @@ -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 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; } diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs index 490579b1e..35a800671 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs @@ -26,10 +26,6 @@ namespace Elsa.Services.Models public IWorkflowBlueprint WorkflowBlueprint { get; } public IServiceProvider ServiceProvider { get; } public WorkflowInstance WorkflowInstance { get; } - - public HashSet BlockingActivities { get; } = - new HashSet(new BlockingActivityEqualityComparer()); - public bool HasScheduledActivities => WorkflowInstance.ScheduledActivities.Any(); public IWorkflowFault? WorkflowFault { get; private set; } public bool IsFirstPass { get; private set; } diff --git a/src/core/Elsa.Core/Activities/ControlFlow/Complete/Complete.cs b/src/core/Elsa.Core/Activities/ControlFlow/Complete/Complete.cs index a82e9c899..a9c1d1885 100644 --- a/src/core/Elsa.Core/Activities/ControlFlow/Complete/Complete.cs +++ b/src/core/Elsa.Core/Activities/ControlFlow/Complete/Complete.cs @@ -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); diff --git a/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs b/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs index 07adbc490..330c969c5 100644 --- a/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs +++ b/src/core/Elsa.Core/Activities/ControlFlow/Join/Join.cs @@ -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) diff --git a/src/core/Elsa.Core/Builders/ActivityBuilder.cs b/src/core/Elsa.Core/Builders/ActivityBuilder.cs index 34259692a..253255353 100644 --- a/src/core/Elsa.Core/Builders/ActivityBuilder.cs +++ b/src/core/Elsa.Core/Builders/ActivityBuilder.cs @@ -56,7 +56,6 @@ namespace Elsa.Builders public Func> 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; diff --git a/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs b/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs index 37c5a6e37..dd22ed5f4 100644 --- a/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs +++ b/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs @@ -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, diff --git a/src/core/Elsa.Core/Services/WorkflowHost.cs b/src/core/Elsa.Core/Services/WorkflowHost.cs index 11ef3d0ec..e706bbfb7 100644 --- a/src/core/Elsa.Core/Services/WorkflowHost.cs +++ b/src/core/Elsa.Core/Services/WorkflowHost.cs @@ -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); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowScheduler.cs b/src/core/Elsa.Core/Services/WorkflowScheduler.cs index 6dd6bda7d..7e3f4c5fd 100644 --- a/src/core/Elsa.Core/Services/WorkflowScheduler.cs +++ b/src/core/Elsa.Core/Services/WorkflowScheduler.cs @@ -151,9 +151,7 @@ namespace Elsa.Services else predicate = x => x.ActivityType == activityType; - var query = _workflowInstanceManager.Query() - .Where(predicate); - + var query = _workflowInstanceManager.Query().Where(predicate); var workflowInstances = await query.ListAsync(); var tuples = workflowInstances.GetBlockingActivities();