From d9d567a48950682bb40431a6d17dbf7a844dcd23 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 11 Oct 2020 22:02:14 +0200 Subject: [PATCH] Incremental work on blueprints --- .../Handlers/TriggerRequestHandler.cs | 6 +- .../ActivityResults/FaultResult.cs | 2 +- .../ActivityResults/OutcomeResult.cs | 2 +- .../ActivityResults/SuspendResult.cs | 2 +- .../Builders/IActivityBuilder.cs | 10 +- .../Elsa.Abstractions/Builders/IBuilder.cs | 2 +- .../Builders/IOutcomeBuilder.cs | 2 +- .../Builders/IWorkflowBuilder.cs | 13 +- .../Extensions/ActivityResolverExtensions.cs | 19 -- .../WorkflowDefinitionExtensions.cs | 29 +++ .../WorkflowExecutionContextExtensions.cs | 22 +- .../Extensions/WorkflowExtensions.cs | 30 +-- .../Messaging/Domain/ActivityNotification.cs | 2 +- .../Models/ActivityDefinition.cs | 28 +-- .../Serialization/IActivitySerializer.cs | 11 + .../Elsa.Abstractions/Services/Activity.cs | 14 +- .../Services/ActivityPropertyProviders.cs | 73 +++++++ .../Elsa.Abstractions/Services/IActivity.cs | 8 +- .../Services/IActivityActivator.cs | 15 ++ .../Services/IActivityPropertyProviders.cs | 17 ++ .../Services/IActivityResolver.cs | 13 -- ...rkflowActivator.cs => IWorkflowFactory.cs} | 6 +- .../Services/IWorkflowHost.cs | 32 +-- .../Services/IWorkflowProvider.cs | 2 +- .../Services/IWorkflowRegistry.cs | 4 +- .../Services/IWorkflowSchedulerQueue.cs | 4 +- .../Services/Models/ActivityBlueprint.cs | 26 +++ .../Models/ActivityExecutionContext.cs | 32 +-- .../Services/Models/Connection.cs | 14 +- .../Services/Models/ExecutionLogEntry.cs | 2 +- .../Services/Models/IActivityBlueprint.cs | 10 + .../Services/Models/IConnection.cs | 8 + .../Services/Models/IEndpoint.cs | 7 + .../Services/Models/IExecutionLogEntry.cs | 10 + .../Services/Models/IScheduledActivity.cs | 10 + .../Services/Models/ISourceEndpoint.cs | 7 + .../Services/Models/ITargetEndpoint.cs | 6 + .../Services/Models/IWorkflowBlueprint.cs | 22 ++ .../Services/Models/IWorkflowFault.cs | 10 + .../Services/Models/ScheduledActivity.cs | 15 +- .../Services/Models/SourceEndpoint.cs | 12 +- .../Services/Models/TargetEndpoint.cs | 6 +- .../{Workflow.cs => WorkflowBlueprint.cs} | 45 ++-- .../Models/WorkflowExecutionContext.cs | 120 +++++------ .../Services/Models/WorkflowFault.cs | 2 +- .../Elsa.Core/Builders/ActivityBuilder.cs | 42 +++- src/core/Elsa.Core/Builders/OutcomeBuilder.cs | 2 +- .../Elsa.Core/Builders/WorkflowBuilder.cs | 138 ++++++------ .../Data/Services/DatabaseWorkflowProvider.cs | 14 +- .../ElsaServiceCollectionExtensions.cs | 5 +- .../Extensions/WorkflowRegistryExtensions.cs | 6 +- .../Serialization/ActivitySerializer.cs | 13 ++ ...tivityResolver.cs => ActivityActivator.cs} | 130 ++++++------ .../Elsa.Core/Services/WorkflowActivator.cs | 39 ---- .../Elsa.Core/Services/WorkflowFactory.cs | 47 +++++ src/core/Elsa.Core/Services/WorkflowHost.cs | 197 ++++++++---------- .../Elsa.Core/Services/WorkflowRegistry.cs | 6 +- .../Elsa.Core/Services/WorkflowScheduler.cs | 26 +-- .../Services/WorkflowSchedulerQueue.cs | 10 +- .../WorkflowProviders/CodeWorkflowProvider.cs | 4 +- .../Elsa.Samples.Serialization/Program.cs | 2 +- .../Program.cs | 6 +- .../Mapping/ActivityStateResolver.cs | 6 +- src/server/Elsa.Server.GraphQL/Query.cs | 8 +- .../Elsa.Core.UnitTests/WorkflowHostTests.cs | 10 +- 65 files changed, 821 insertions(+), 622 deletions(-) delete mode 100644 src/core/Elsa.Abstractions/Extensions/ActivityResolverExtensions.cs create mode 100644 src/core/Elsa.Abstractions/Extensions/WorkflowDefinitionExtensions.cs create mode 100644 src/core/Elsa.Abstractions/Serialization/IActivitySerializer.cs create mode 100644 src/core/Elsa.Abstractions/Services/ActivityPropertyProviders.cs create mode 100644 src/core/Elsa.Abstractions/Services/IActivityActivator.cs create mode 100644 src/core/Elsa.Abstractions/Services/IActivityPropertyProviders.cs delete mode 100644 src/core/Elsa.Abstractions/Services/IActivityResolver.cs rename src/core/Elsa.Abstractions/Services/{IWorkflowActivator.cs => IWorkflowFactory.cs} (63%) create mode 100644 src/core/Elsa.Abstractions/Services/Models/ActivityBlueprint.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/IActivityBlueprint.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/IConnection.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/IEndpoint.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/IExecutionLogEntry.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/IScheduledActivity.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/ISourceEndpoint.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/ITargetEndpoint.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs create mode 100644 src/core/Elsa.Abstractions/Services/Models/IWorkflowFault.cs rename src/core/Elsa.Abstractions/Services/Models/{Workflow.cs => WorkflowBlueprint.cs} (56%) create mode 100644 src/core/Elsa.Core/Serialization/ActivitySerializer.cs rename src/core/Elsa.Core/Services/{ActivityResolver.cs => ActivityActivator.cs} (64%) delete mode 100644 src/core/Elsa.Core/Services/WorkflowActivator.cs create mode 100644 src/core/Elsa.Core/Services/WorkflowFactory.cs diff --git a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs index 2a30038fc..228432064 100644 --- a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs +++ b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs @@ -74,8 +74,8 @@ namespace Elsa.Activities.Http.RequestHandlers.Handlers : new EmptyResult(); } - private IEnumerable<(Workflow Workflow, ReceiveHttpRequest Activity)> Filter( - IEnumerable<(Workflow Workflow, ReceiveHttpRequest Activity)> items, + private IEnumerable<(WorkflowBlueprint Workflow, ReceiveHttpRequest Activity)> Filter( + IEnumerable<(WorkflowBlueprint Workflow, ReceiveHttpRequest Activity)> items, PathString path, string method) => items.Where(x => IsMatch(x.Activity, path, method)); @@ -88,7 +88,7 @@ namespace Elsa.Activities.Http.RequestHandlers.Handlers } private async Task InvokeWorkflowsToStartAsync( - IEnumerable<(Workflow Workflow, ReceiveHttpRequest Activity)> items) + IEnumerable<(WorkflowBlueprint Workflow, ReceiveHttpRequest Activity)> items) { foreach (var (workflow, activity) in items) { diff --git a/src/core/Elsa.Abstractions/ActivityResults/FaultResult.cs b/src/core/Elsa.Abstractions/ActivityResults/FaultResult.cs index ef38c9fa0..5f3ff4991 100644 --- a/src/core/Elsa.Abstractions/ActivityResults/FaultResult.cs +++ b/src/core/Elsa.Abstractions/ActivityResults/FaultResult.cs @@ -9,6 +9,6 @@ namespace Elsa.ActivityResults public LocalizedString Message { get; } protected override void Execute(ActivityExecutionContext activityExecutionContext) => - activityExecutionContext.WorkflowExecutionContext.Fault(activityExecutionContext.Activity, Message); + activityExecutionContext.WorkflowExecutionContext.Fault(activityExecutionContext.ActivityDefinition, Message); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/ActivityResults/OutcomeResult.cs b/src/core/Elsa.Abstractions/ActivityResults/OutcomeResult.cs index 6b5f62d44..7527b4fe3 100644 --- a/src/core/Elsa.Abstractions/ActivityResults/OutcomeResult.cs +++ b/src/core/Elsa.Abstractions/ActivityResults/OutcomeResult.cs @@ -30,7 +30,7 @@ namespace Elsa.ActivityResults activityExecutionContext.Outcomes = Outcomes.ToList(); var workflowExecutionContext = activityExecutionContext.WorkflowExecutionContext; - var nextActivities = GetNextActivities(workflowExecutionContext, activityExecutionContext.Activity, Outcomes).ToList(); + var nextActivities = GetNextActivities(workflowExecutionContext, activityExecutionContext.ActivityDefinition, Outcomes).ToList(); workflowExecutionContext.ScheduleActivities(nextActivities, Output); } diff --git a/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs b/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs index cdae0bd91..39bcfc603 100644 --- a/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs +++ b/src/core/Elsa.Abstractions/ActivityResults/SuspendResult.cs @@ -6,7 +6,7 @@ namespace Elsa.ActivityResults { protected override void Execute(ActivityExecutionContext activityExecutionContext) { - activityExecutionContext.WorkflowExecutionContext.BlockingActivities.Add(activityExecutionContext.Activity); + activityExecutionContext.WorkflowExecutionContext.BlockingActivities.Add(activityExecutionContext.ActivityDefinition); activityExecutionContext.WorkflowExecutionContext.Suspend(); } } diff --git a/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs b/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs index 51fff3741..e22f9cf95 100644 --- a/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs +++ b/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs @@ -1,5 +1,7 @@ using System; using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; using Elsa.Services; using Elsa.Services.Models; @@ -8,12 +10,14 @@ namespace Elsa.Builders public interface IActivityBuilder : IBuilder { IWorkflowBuilder WorkflowBuilder { get; } - IActivity Activity { get; } + public Type ActivityType { get; } + string? ActivityId { get; } IDictionary? PropertyValueProviders { get; } IActivityBuilder Add(Action>? setup = default) where T : class, IActivity; IOutcomeBuilder When(string outcome); IActivityBuilder Then(IActivityBuilder targetActivity); - IActivity BuildActivity(); - Workflow Build(); + IActivityBuilder WithId(string id); + Func> BuildActivityAsync(); + IWorkflowBlueprint Build(); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Builders/IBuilder.cs b/src/core/Elsa.Abstractions/Builders/IBuilder.cs index 56cb89c1e..c46ee81bb 100644 --- a/src/core/Elsa.Abstractions/Builders/IBuilder.cs +++ b/src/core/Elsa.Abstractions/Builders/IBuilder.cs @@ -7,6 +7,6 @@ namespace Elsa.Builders { IActivityBuilder Then(Action>? setup = default, Action? branch = default) where T : class, IActivity; IActivityBuilder Then(Action setup, Action? branch = default) where T : class, IActivity; - IActivityBuilder Then(T activity, Action? branch = default) where T : class, IActivity; + IActivityBuilder Then(Action? branch = default) where T : class, IActivity; } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs b/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs index 4932277e5..e863729ac 100644 --- a/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs +++ b/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs @@ -7,6 +7,6 @@ namespace Elsa.Builders IWorkflowBuilder WorkflowBuilder { get; } IActivityBuilder Source { get; } string? Outcome { get; } - Workflow Build(); + WorkflowBlueprint Build(); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Builders/IWorkflowBuilder.cs b/src/core/Elsa.Abstractions/Builders/IWorkflowBuilder.cs index 7c6510ac5..0202018be 100644 --- a/src/core/Elsa.Abstractions/Builders/IWorkflowBuilder.cs +++ b/src/core/Elsa.Abstractions/Builders/IWorkflowBuilder.cs @@ -24,7 +24,7 @@ namespace Elsa.Builders IWorkflowBuilder WithDeleteCompletedInstances(bool value); IWorkflowBuilder WithPersistenceBehavior(WorkflowPersistenceBehavior value); - IActivityBuilder New(T activity, + IActivityBuilder New( Action? branch = default, IDictionary? propertyValueProviders = default) where T : class, IActivity; @@ -36,7 +36,7 @@ namespace Elsa.Builders Action>? setup = default, Action? branch = default) where T : class, IActivity; - IActivityBuilder StartWith(T activity, Action? branch = default) + IActivityBuilder StartWith(Action? branch = default) where T : class, IActivity; IActivityBuilder Add( @@ -48,7 +48,6 @@ namespace Elsa.Builders Action? branch = default) where T : class, IActivity; IActivityBuilder Add( - T activity, Action? branch = default, IDictionary? propertyValueProviders = default) where T : class, IActivity; @@ -62,9 +61,9 @@ namespace Elsa.Builders Func target, string outcome = OutcomeNames.Done); - Workflow Build(); - Workflow Build(IWorkflow workflow); - Workflow Build(Type workflowType); - Workflow Build() where T : IWorkflow; + IWorkflowBlueprint Build(); + IWorkflowBlueprint Build(IWorkflow workflow); + IWorkflowBlueprint Build(Type workflowType); + IWorkflowBlueprint Build() where T : IWorkflow; } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Extensions/ActivityResolverExtensions.cs b/src/core/Elsa.Abstractions/Extensions/ActivityResolverExtensions.cs deleted file mode 100644 index 327db6764..000000000 --- a/src/core/Elsa.Abstractions/Extensions/ActivityResolverExtensions.cs +++ /dev/null @@ -1,19 +0,0 @@ -using Elsa.Models; -using Elsa.Services; - -namespace Elsa -{ - public static class ActivityResolverExtensions - { - public static IActivity ResolveActivity(this IActivityResolver activityResolver, ActivityDefinition activityDefinition) - { - var activity = activityResolver.ResolveActivity(activityDefinition.Type); - activity.Description = activityDefinition.Description; - activity.Id = activityDefinition.Id; - activity.Name = activityDefinition.Name; - activity.DisplayName = activityDefinition.DisplayName; - activity.PersistWorkflow = activityDefinition.PersistWorkflow; - return activity; - } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowDefinitionExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowDefinitionExtensions.cs new file mode 100644 index 000000000..50c4fb99d --- /dev/null +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowDefinitionExtensions.cs @@ -0,0 +1,29 @@ +using System.Collections.Generic; +using System.Linq; +using Elsa.Models; + +namespace Elsa +{ + public static class WorkflowDefinitionExtensions + { + public static ActivityDefinition + GetActivityById(this WorkflowDefinition workflowDefinition, string activityId) => + workflowDefinition.Activities.First(x => x.Id == activityId); + + public static IEnumerable GetStartActivities(this WorkflowDefinition workflowDefinition) + { + var targetActivities = workflowDefinition.Connections + .Select(x => x.TargetActivityId) + .Where(x => x != null) + .Distinct() + .ToLookup(x => x); + + var query = + from activity in workflowDefinition.Activities + where !targetActivities.Contains(activity.Id) + select activity; + + return query; + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs index 43eb8b3a2..b1bb90e32 100644 --- a/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowExecutionContextExtensions.cs @@ -7,17 +7,17 @@ namespace Elsa { public static class WorkflowExecutionContextExtensions { - public static IEnumerable GetStartActivities(this WorkflowExecutionContext workflowExecutionContext) - { - var targetActivities = workflowExecutionContext.Connections.Select(x => x.Target.Activity).Distinct().ToLookup(x => x); - - var query = - from activity in workflowExecutionContext.Activities - where !targetActivities.Contains(activity) - select activity; - - return query; - } + // public static IEnumerable GetStartActivities(this WorkflowExecutionContext workflowExecutionContext) + // { + // var targetActivities = workflowExecutionContext.Connections.Select(x => x.Target.Activity).Distinct().ToLookup(x => x); + // + // var query = + // from activity in workflowExecutionContext.Activities + // where !targetActivities.Contains(activity) + // select activity; + // + // return query; + // } public static IEnumerable GetInboundConnections(this WorkflowExecutionContext workflowExecutionContext, IActivity activity) => workflowExecutionContext.Connections.Where(x => x.Target.Activity == activity); public static IEnumerable GetOutboundConnections(this WorkflowExecutionContext workflowExecutionContext, IActivity activity) => workflowExecutionContext.Connections.Where(x => x.Source.Activity == activity); diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowExtensions.cs index f1c639268..4f0cb4b3f 100644 --- a/src/core/Elsa.Abstractions/Extensions/WorkflowExtensions.cs +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowExtensions.cs @@ -8,10 +8,10 @@ namespace Elsa { public static class WorkflowExtensions { - public static IEnumerable WithVersion(this IEnumerable query, VersionOptions version) + public static IEnumerable WithVersion(this IEnumerable query, VersionOptions version) => query.AsQueryable().WithVersion(version); - public static IQueryable WithVersion(this IQueryable query, VersionOptions version) + public static IQueryable WithVersion(this IQueryable query, VersionOptions version) { if (version.IsDraft) query = query.Where(x => !x.IsPublished); @@ -31,47 +31,47 @@ namespace Elsa return query.OrderByDescending(x => x.Version); } - public static IEnumerable GetStartActivities(this Workflow workflow) + public static IEnumerable GetStartActivities(this WorkflowBlueprint workflowBlueprint) { - var targetActivityIds = workflow.Connections.Select(x => x.Target.Activity.Id).Distinct().ToLookup(x => x); + var targetActivityIds = workflowBlueprint.Connections.Select(x => x.Target.Activity.Id).Distinct().ToLookup(x => x); var query = - from activity in workflow.Activities + from activity in workflowBlueprint.Activities where !targetActivityIds.Contains(activity.Id) select activity; return query; } - public static IActivity GetActivity(this Workflow workflow, string id) => workflow.Activities.FirstOrDefault(x => x.Id == id); + public static IActivity GetActivity(this WorkflowBlueprint workflowBlueprint, string id) => workflowBlueprint.Activities.FirstOrDefault(x => x.Id == id); - public static IEnumerable GetInboundConnections(this Workflow workflow, string activityId) + public static IEnumerable GetInboundConnections(this WorkflowBlueprint workflowBlueprint, string activityId) { - return workflow.Connections.Where(x => x.Target.Activity.Id == activityId).ToList(); + return workflowBlueprint.Connections.Where(x => x.Target.Activity.Id == activityId).ToList(); } - public static IEnumerable GetOutboundConnections(this Workflow workflow, string activityId) + public static IEnumerable GetOutboundConnections(this WorkflowBlueprint workflowBlueprint, string activityId) { - return workflow.Connections.Where(x => x.Source.Activity.Id == activityId).ToList(); + return workflowBlueprint.Connections.Where(x => x.Source.Activity.Id == activityId).ToList(); } /// /// Returns the full path of incoming activities. /// - public static IEnumerable GetInboundActivityPath(this Workflow workflow, string activityId) + public static IEnumerable GetInboundActivityPath(this WorkflowBlueprint workflowBlueprint, string activityId) { var inspectedActivityIDs = new HashSet(); - return workflow.GetInboundActivityPathInternal(activityId, activityId, inspectedActivityIDs) + return workflowBlueprint.GetInboundActivityPathInternal(activityId, activityId, inspectedActivityIDs) .Distinct().ToList(); } - private static IEnumerable GetInboundActivityPathInternal(this Workflow workflowInstance, + private static IEnumerable GetInboundActivityPathInternal(this WorkflowBlueprint workflowBlueprintBlueprintInstance, string activityId, string startingPointActivityId, HashSet inspectedActivityIDs) { - foreach (var connection in workflowInstance.GetInboundConnections(activityId)) + foreach (var connection in workflowBlueprintBlueprintInstance.GetInboundConnections(activityId)) { // Circuit breaker: Detect workflows that implement repeating flows to prevent an infinite loop here. if (inspectedActivityIDs.Contains(connection.Source.Activity.Id)) @@ -79,7 +79,7 @@ namespace Elsa yield return connection.Source.Activity.Id; - foreach (var parentActivityId in workflowInstance.GetInboundActivityPathInternal(connection.Source.Activity.Id, startingPointActivityId, inspectedActivityIDs) + foreach (var parentActivityId in workflowBlueprintBlueprintInstance.GetInboundActivityPathInternal(connection.Source.Activity.Id, startingPointActivityId, inspectedActivityIDs) .Distinct()) { inspectedActivityIDs.Add(parentActivityId); diff --git a/src/core/Elsa.Abstractions/Messaging/Domain/ActivityNotification.cs b/src/core/Elsa.Abstractions/Messaging/Domain/ActivityNotification.cs index 9100a794a..a1fe7fe2a 100644 --- a/src/core/Elsa.Abstractions/Messaging/Domain/ActivityNotification.cs +++ b/src/core/Elsa.Abstractions/Messaging/Domain/ActivityNotification.cs @@ -13,6 +13,6 @@ namespace Elsa.Messaging.Domain public ActivityExecutionContext ActivityExecutionContext { get; } public WorkflowExecutionContext WorkflowExecutionContext => ActivityExecutionContext.WorkflowExecutionContext; - public IActivity Activity => ActivityExecutionContext.Activity; + public IActivity Activity => ActivityExecutionContext.ActivityDefinition; } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/ActivityDefinition.cs b/src/core/Elsa.Abstractions/Models/ActivityDefinition.cs index 800b7bc6f..3dd2e0001 100644 --- a/src/core/Elsa.Abstractions/Models/ActivityDefinition.cs +++ b/src/core/Elsa.Abstractions/Models/ActivityDefinition.cs @@ -1,37 +1,17 @@ -using Elsa.Services; +using Newtonsoft.Json.Linq; namespace Elsa.Models { public class ActivityDefinition { - public static ActivityDefinition FromActivity(IActivity activity) - { - return new ActivityDefinition - { - Id = activity.Id, - Type = activity.Type, - //State = activity.State, - Name = activity.Name, - DisplayName = activity.DisplayName - }; - } - - public string Id { get; set; } - public string Type { get; set; } + public string Id { get; set; } = default!; + public string Type { get; set; } = default!; public string? Name { get; set; } public string? DisplayName { get; set; } public string? Description { get; set; } public int? Left { get; set; } public int? Top { get; set; } public bool PersistWorkflow { get; set; } - public Variables? State { get; set; } - } - - public class ActivityDefinition : ActivityDefinition where T : IActivity - { - public ActivityDefinition() - { - Type = typeof(T).Name; - } + public JObject Data { get; set; } = new JObject(); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Serialization/IActivitySerializer.cs b/src/core/Elsa.Abstractions/Serialization/IActivitySerializer.cs new file mode 100644 index 000000000..c9c11be9f --- /dev/null +++ b/src/core/Elsa.Abstractions/Serialization/IActivitySerializer.cs @@ -0,0 +1,11 @@ +using Elsa.Services; +using Newtonsoft.Json.Linq; + +namespace Elsa.Serialization +{ + public interface IActivitySerializer + { + JObject Serialize(T activity) where T : IActivity; + T Deserialize(JObject data) where T : IActivity; + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Activity.cs b/src/core/Elsa.Abstractions/Services/Activity.cs index 649326ce3..94be90adf 100644 --- a/src/core/Elsa.Abstractions/Services/Activity.cs +++ b/src/core/Elsa.Abstractions/Services/Activity.cs @@ -18,16 +18,16 @@ namespace Elsa.Services public string? Description{ get; set; } public bool PersistWorkflow { get; set; } - public Task CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(context, cancellationToken); - public Task ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnExecuteAsync(context, cancellationToken); + public ValueTask CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(context, cancellationToken); + public ValueTask ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnExecuteAsync(context, cancellationToken); - public Task ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnResumeAsync(context, cancellationToken); + public ValueTask ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnResumeAsync(context, cancellationToken); protected virtual bool OnCanExecute(ActivityExecutionContext context) => OnCanExecute(); protected virtual bool OnCanExecute() => true; - protected virtual Task OnCanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(cancellationToken); - protected virtual Task OnCanExecuteAsync(CancellationToken cancellationToken) => Task.FromResult(OnCanExecute()); - protected virtual Task OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => Task.FromResult(OnExecute(context)); - protected virtual Task OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => Task.FromResult(OnResume(context)); + protected virtual ValueTask OnCanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(cancellationToken); + protected virtual ValueTask OnCanExecuteAsync(CancellationToken cancellationToken) => new ValueTask(OnCanExecute()); + protected virtual ValueTask OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => new ValueTask(OnExecute(context)); + protected virtual ValueTask OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => new ValueTask(OnResume(context)); protected virtual IActivityExecutionResult OnExecute(ActivityExecutionContext context) => OnExecute(); protected virtual IActivityExecutionResult OnExecute() => Done(); protected virtual IActivityExecutionResult OnResume(ActivityExecutionContext context) => OnResume(); diff --git a/src/core/Elsa.Abstractions/Services/ActivityPropertyProviders.cs b/src/core/Elsa.Abstractions/Services/ActivityPropertyProviders.cs new file mode 100644 index 000000000..83021f0fc --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/ActivityPropertyProviders.cs @@ -0,0 +1,73 @@ +using System.Collections.Generic; +using System.Linq; +using System.Reflection; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Attributes; +using Elsa.Models; +using Elsa.Services.Models; + +namespace Elsa.Services +{ + public class ActivityPropertyProviders : IActivityPropertyProviders + { + private readonly IDictionary> _providers = + new Dictionary>(); + + public ActivityPropertyProviders() + { + } + + public ActivityPropertyProviders(IDictionary> providers) + { + _providers = providers; + } + + public void AddProvider(string activityId, string propertyName, IActivityPropertyValueProvider provider) + { + if (!_providers.TryGetValue(activityId, out var properties)) + { + properties = new Dictionary(); + _providers.Add(activityId, properties); + } + + properties[propertyName] = provider; + } + + public IDictionary? GetProviders(string activityId) => + _providers.TryGetValue(activityId, out var properties) ? properties : null; + + public IActivityPropertyValueProvider? GetProvider(string activityId, string propertyName) + { + if (_providers.TryGetValue(activityId, out var properties)) + if (properties.TryGetValue(propertyName, out var provider)) + return provider; + + return null; + } + + public async ValueTask SetActivityPropertiesAsync( + IActivity activity, + ActivityExecutionContext activityExecutionContext, + CancellationToken cancellationToken = default) + { + var properties = activity.GetType().GetProperties().Where(IsActivityProperty).ToList(); + var providers = GetProviders(activity.Id); + + if (providers == null) + return; + + foreach (var property in properties) + { + if (!providers.TryGetValue(property.Name, out var provider)) + continue; + + var value = await provider.GetValueAsync(activityExecutionContext, cancellationToken); + property.SetValue(activity, value); + } + } + + private bool IsActivityProperty(PropertyInfo property) => + property.GetCustomAttribute() != null; + } +} \ 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 4114febd9..ff146a673 100644 --- a/src/core/Elsa.Abstractions/Services/IActivity.cs +++ b/src/core/Elsa.Abstractions/Services/IActivity.cs @@ -13,7 +13,7 @@ namespace Elsa.Services string Type { get; } /// - /// Unique identifier of this activity. + /// Unique identifier of this activity within the workflow. /// string Id { get; set; } @@ -45,16 +45,16 @@ namespace Elsa.Services /// /// Returns a value of whether the specified activity can execute. /// - Task CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default); + ValueTask CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default); /// /// Executes the specified activity. /// - Task ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default); + ValueTask ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default); /// /// Resumes the specified activity. /// - Task ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default); + ValueTask ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IActivityActivator.cs b/src/core/Elsa.Abstractions/Services/IActivityActivator.cs new file mode 100644 index 000000000..b3eca2421 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/IActivityActivator.cs @@ -0,0 +1,15 @@ +using System; +using System.Collections.Generic; +using Elsa.Models; + +namespace Elsa.Services +{ + public interface IActivityActivator + { + IActivity ActivateActivity(string activityTypeName, Action? setup = default); + T ActivateActivity(Action? configure = default) where T : class, IActivity; + IActivity ActivateActivity(ActivityDefinition activityDefinition); + IEnumerable GetActivityTypes(); + Type? GetActivityType(string activityTypeName); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IActivityPropertyProviders.cs b/src/core/Elsa.Abstractions/Services/IActivityPropertyProviders.cs new file mode 100644 index 000000000..14d884817 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/IActivityPropertyProviders.cs @@ -0,0 +1,17 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services.Models; + +namespace Elsa.Services +{ + public interface IActivityPropertyProviders + { + void AddProvider(string activityId, string propertyName, IActivityPropertyValueProvider provider); + IActivityPropertyValueProvider? GetProvider(string activityId, string propertyName); + + ValueTask SetActivityPropertiesAsync( + IActivity activity, + ActivityExecutionContext activityExecutionContext, + CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IActivityResolver.cs b/src/core/Elsa.Abstractions/Services/IActivityResolver.cs deleted file mode 100644 index 9f1da1985..000000000 --- a/src/core/Elsa.Abstractions/Services/IActivityResolver.cs +++ /dev/null @@ -1,13 +0,0 @@ -using System; -using System.Collections.Generic; - -namespace Elsa.Services -{ - public interface IActivityResolver - { - IActivity ResolveActivity(string activityTypeName, Action? setup = default); - T ResolveActivity(Action? configure = default) where T : class, IActivity; - IEnumerable GetActivityTypes(); - Type? GetActivityType(string activityTypeName); - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowActivator.cs b/src/core/Elsa.Abstractions/Services/IWorkflowFactory.cs similarity index 63% rename from src/core/Elsa.Abstractions/Services/IWorkflowActivator.cs rename to src/core/Elsa.Abstractions/Services/IWorkflowFactory.cs index d43472288..ae9920426 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowActivator.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowFactory.cs @@ -5,10 +5,10 @@ using Elsa.Services.Models; namespace Elsa.Services { - public interface IWorkflowActivator + public interface IWorkflowFactory { - Task ActivateAsync( - Workflow workflow, + Task InstantiateAsync( + WorkflowDefinition workflowDefinition, string? correlationId = default, CancellationToken cancellationToken = default); } diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs b/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs index 620cc5368..43a079c20 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs @@ -7,20 +7,24 @@ namespace Elsa.Services { public interface IWorkflowHost { - Task RunWorkflowAsync(Workflow workflow, string? activityId = default, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default); + ValueTask RunWorkflowInstanceAsync( + WorkflowInstance workflowInstance, + string? activityId = default, + object? input = default, + CancellationToken cancellationToken = default); + + ValueTask RunWorkflowInstanceAsync( + WorkflowDefinition workflowDefinition, + WorkflowInstance workflowInstance, + string? activityId = default, + object? input = default, + CancellationToken cancellationToken = default); - Task RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default); - - Task RunWorkflowInstanceAsync(WorkflowInstance workflowInstance, string? activityId = default, object? input = default, CancellationToken cancellationToken = default); - - /// - /// Run a registered workflow by its ID. - /// - Task RunWorkflowDefinitionAsync(string workflowDefinitionId, string? activityId = default, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default); - - // /// - // /// Resume a workflow instance. - // /// - // Task ResumeAsync(string workflowInstanceId, string activityId, object? input = default, CancellationToken cancellationToken = default); + ValueTask RunWorkflowDefinitionAsync( + WorkflowDefinition workflowDefinition, + string? activityId = default, + object? input = default, + string? correlationId = default, + CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowProvider.cs b/src/core/Elsa.Abstractions/Services/IWorkflowProvider.cs index 4cef03147..6316f9748 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowProvider.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowProvider.cs @@ -10,6 +10,6 @@ namespace Elsa.Services /// public interface IWorkflowProvider { - Task> GetWorkflowsAsync(CancellationToken cancellationToken); + Task> GetWorkflowsAsync(CancellationToken cancellationToken); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowRegistry.cs b/src/core/Elsa.Abstractions/Services/IWorkflowRegistry.cs index 3223fef05..d4b33e114 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowRegistry.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowRegistry.cs @@ -8,7 +8,7 @@ namespace Elsa.Services { public interface IWorkflowRegistry { - Task> GetWorkflowsAsync(CancellationToken cancellationToken = default); - Task GetWorkflowAsync(string id, VersionOptions version, CancellationToken cancellationToken = default); + Task> GetWorkflowsAsync(CancellationToken cancellationToken = default); + Task GetWorkflowAsync(string id, VersionOptions version, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowSchedulerQueue.cs b/src/core/Elsa.Abstractions/Services/IWorkflowSchedulerQueue.cs index ca0941c3d..bb8d4e447 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowSchedulerQueue.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowSchedulerQueue.cs @@ -4,7 +4,7 @@ namespace Elsa.Services { public interface IWorkflowSchedulerQueue { - void Enqueue(Workflow workflow, IActivity activity, object? input, string? correlationId); - (Workflow Workflow, IActivity Activity, object? Input, string? CorrelationId)? Dequeue(string workflowDefinitionId, string activityId); + void Enqueue(WorkflowBlueprint workflowBlueprint, IActivity activity, object? input, string? correlationId); + (WorkflowBlueprint Workflow, IActivity Activity, object? Input, string? CorrelationId)? Dequeue(string workflowDefinitionId, string activityId); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/ActivityBlueprint.cs b/src/core/Elsa.Abstractions/Services/Models/ActivityBlueprint.cs new file mode 100644 index 000000000..a77fe5d04 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/ActivityBlueprint.cs @@ -0,0 +1,26 @@ +using System; +using Elsa.ActivityResults; + +namespace Elsa.Services.Models +{ + public class ActivityBlueprint : IActivityBlueprint + { + public ActivityBlueprint() + { + } + + public ActivityBlueprint(Func createActivity) + { + CreateActivity = createActivity; + } + + public ActivityBlueprint(string id, Func createActivity) + { + Id = id; + CreateActivity = createActivity; + } + + public string Id { get; set; } = default!; + public Func CreateActivity { get; set; } = default!; + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs index e1bca2324..77bb5b19e 100644 --- a/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs @@ -4,6 +4,7 @@ using System.Reflection; using System.Threading; using System.Threading.Tasks; using Elsa.Attributes; +using Elsa.Models; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Services.Models @@ -12,17 +13,17 @@ namespace Elsa.Services.Models { public ActivityExecutionContext( WorkflowExecutionContext workflowExecutionContext, - IActivity activity, + ActivityDefinition activityDefinition, object? input = null) { WorkflowExecutionContext = workflowExecutionContext; - Activity = activity; + ActivityDefinition = activityDefinition; Input = input; Outcomes = new List(0); } public WorkflowExecutionContext WorkflowExecutionContext { get; } - public IActivity Activity { get; } + public ActivityDefinition ActivityDefinition { get; } public object? Input { get; } public object? Output { get; set; } public IReadOnlyCollection Outcomes { get; set; } @@ -32,24 +33,11 @@ namespace Elsa.Services.Models public T GetVariable(string name) => WorkflowExecutionContext.GetVariable(name); public T GetService() => WorkflowExecutionContext.ServiceProvider.GetService(); - public async ValueTask SetActivityPropertiesAsync(CancellationToken cancellationToken = default) - { - var properties = Activity.GetType().GetProperties().Where(IsActivityProperty).ToList(); - var activityPropertyValueProviders = WorkflowExecutionContext.ActivityPropertyValueProviders; - var propertyValueProvider = activityPropertyValueProviders[Activity.Id]; - - foreach (var property in properties) - { - if(propertyValueProvider == null || !propertyValueProvider.ContainsKey(property.Name)) - continue; - - var provider = propertyValueProvider[property.Name]; - var value = await provider.GetValueAsync( this, cancellationToken); - property.SetValue(Activity, value); - } - } - - private bool IsActivityProperty(PropertyInfo property) => - property.GetCustomAttribute() != null; + public async ValueTask SetActivityPropertiesAsync(IActivity activity, + CancellationToken cancellationToken = default) => + await WorkflowExecutionContext.ActivityPropertyProviders.SetActivityPropertiesAsync( + activity, + this, + cancellationToken); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/Connection.cs b/src/core/Elsa.Abstractions/Services/Models/Connection.cs index 375642923..40cf8eabf 100644 --- a/src/core/Elsa.Abstractions/Services/Models/Connection.cs +++ b/src/core/Elsa.Abstractions/Services/Models/Connection.cs @@ -1,23 +1,19 @@ namespace Elsa.Services.Models { - public class Connection + public class Connection : IConnection { - public Connection() - { - } - - public Connection(IActivity sourceActivity, IActivity targetActivity, string sourceOutcome = OutcomeNames.Done) + public Connection(IActivity sourceActivity, IActivity targetActivity, string sourceOutcome) : this(new SourceEndpoint(sourceActivity, sourceOutcome), new TargetEndpoint(targetActivity)) { } - public Connection(SourceEndpoint source, TargetEndpoint target) + public Connection(ISourceEndpoint source, ITargetEndpoint target) { Source = source; Target = target; } - public SourceEndpoint Source { get; set; } - public TargetEndpoint Target { get; set; } + public ISourceEndpoint Source { get; set; } + public ITargetEndpoint Target { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/ExecutionLogEntry.cs b/src/core/Elsa.Abstractions/Services/Models/ExecutionLogEntry.cs index 9b1505f62..c030e4b8f 100644 --- a/src/core/Elsa.Abstractions/Services/Models/ExecutionLogEntry.cs +++ b/src/core/Elsa.Abstractions/Services/Models/ExecutionLogEntry.cs @@ -2,7 +2,7 @@ using NodaTime; namespace Elsa.Services.Models { - public class ExecutionLogEntry + public class ExecutionLogEntry : IExecutionLogEntry { public ExecutionLogEntry(IActivity activity, Instant timestamp) { diff --git a/src/core/Elsa.Abstractions/Services/Models/IActivityBlueprint.cs b/src/core/Elsa.Abstractions/Services/Models/IActivityBlueprint.cs new file mode 100644 index 000000000..e07afee71 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IActivityBlueprint.cs @@ -0,0 +1,10 @@ +using System; + +namespace Elsa.Services.Models +{ + public interface IActivityBlueprint + { + public string Id { get; } + Func CreateActivity { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IConnection.cs b/src/core/Elsa.Abstractions/Services/Models/IConnection.cs new file mode 100644 index 000000000..3b235f1f7 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IConnection.cs @@ -0,0 +1,8 @@ +namespace Elsa.Services.Models +{ + public interface IConnection + { + ISourceEndpoint Source { get; } + ITargetEndpoint Target { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IEndpoint.cs b/src/core/Elsa.Abstractions/Services/Models/IEndpoint.cs new file mode 100644 index 000000000..3fbb31e8a --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IEndpoint.cs @@ -0,0 +1,7 @@ +namespace Elsa.Services.Models +{ + public interface IEndpoint + { + IActivity Activity { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IExecutionLogEntry.cs b/src/core/Elsa.Abstractions/Services/Models/IExecutionLogEntry.cs new file mode 100644 index 000000000..dcd45b4b1 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IExecutionLogEntry.cs @@ -0,0 +1,10 @@ +using NodaTime; + +namespace Elsa.Services.Models +{ + public interface IExecutionLogEntry + { + IActivity Activity { get; } + Instant Timestamp { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IScheduledActivity.cs b/src/core/Elsa.Abstractions/Services/Models/IScheduledActivity.cs new file mode 100644 index 000000000..4af190865 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IScheduledActivity.cs @@ -0,0 +1,10 @@ +using Elsa.Models; + +namespace Elsa.Services.Models +{ + public interface IScheduledActivity + { + ActivityDefinition ActivityDefinition { get; } + object? Input { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/ISourceEndpoint.cs b/src/core/Elsa.Abstractions/Services/Models/ISourceEndpoint.cs new file mode 100644 index 000000000..657a1adaf --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/ISourceEndpoint.cs @@ -0,0 +1,7 @@ +namespace Elsa.Services.Models +{ + public interface ISourceEndpoint : IEndpoint + { + string Outcome { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/ITargetEndpoint.cs b/src/core/Elsa.Abstractions/Services/Models/ITargetEndpoint.cs new file mode 100644 index 000000000..f29aee83c --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/ITargetEndpoint.cs @@ -0,0 +1,6 @@ +namespace Elsa.Services.Models +{ + public interface ITargetEndpoint : IEndpoint + { + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs b/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs new file mode 100644 index 000000000..a3672af86 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IWorkflowBlueprint.cs @@ -0,0 +1,22 @@ +using System.Collections.Generic; +using Elsa.Models; + +namespace Elsa.Services.Models +{ + public interface IWorkflowBlueprint + { + public string? Name { get; } + public string Id { get; } + public int Version { get; set; } + public bool IsSingleton { get; } + public bool IsEnabled { get; } + public string? Description { get; } + public bool IsPublished { get; } + public bool IsLatest { get; } + public WorkflowPersistenceBehavior PersistenceBehavior { get; } + public bool DeleteCompletedInstances { get; } + public ICollection Activities { get; } + public ICollection Connections { get; } + IActivityPropertyProviders ActivityPropertyProviders { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/IWorkflowFault.cs b/src/core/Elsa.Abstractions/Services/Models/IWorkflowFault.cs new file mode 100644 index 000000000..ab66a44ba --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/IWorkflowFault.cs @@ -0,0 +1,10 @@ +using Microsoft.Extensions.Localization; + +namespace Elsa.Services.Models +{ + public interface IWorkflowFault + { + IActivity? FaultedActivity { get; } + LocalizedString? Message { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/ScheduledActivity.cs b/src/core/Elsa.Abstractions/Services/Models/ScheduledActivity.cs index f41bfb696..263197589 100644 --- a/src/core/Elsa.Abstractions/Services/Models/ScheduledActivity.cs +++ b/src/core/Elsa.Abstractions/Services/Models/ScheduledActivity.cs @@ -1,15 +1,16 @@ -namespace Elsa.Services.Models +using Elsa.Models; + +namespace Elsa.Services.Models { - public class ScheduledActivity + public class ScheduledActivity : IScheduledActivity { - public ScheduledActivity(IActivity activity, object? input = default) + public ScheduledActivity(ActivityDefinition activityDefinition, object? input = default) { - Activity = activity; + ActivityDefinition = activityDefinition; Input = input; } - - - public IActivity Activity { get; } + + public ActivityDefinition ActivityDefinition { get; } public object? Input { get; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/SourceEndpoint.cs b/src/core/Elsa.Abstractions/Services/Models/SourceEndpoint.cs index 64b938016..60db86184 100644 --- a/src/core/Elsa.Abstractions/Services/Models/SourceEndpoint.cs +++ b/src/core/Elsa.Abstractions/Services/Models/SourceEndpoint.cs @@ -1,16 +1,8 @@ namespace Elsa.Services.Models { - public class SourceEndpoint : Endpoint + public class SourceEndpoint : Endpoint, ISourceEndpoint { - public SourceEndpoint() - { - } - - public SourceEndpoint(IActivity activity, string outcome) : base(activity) - { - Outcome = outcome; - } - + public SourceEndpoint(IActivity activity, string outcome) : base(activity) => Outcome = outcome; public string Outcome { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/TargetEndpoint.cs b/src/core/Elsa.Abstractions/Services/Models/TargetEndpoint.cs index 1814713fb..ab838e119 100644 --- a/src/core/Elsa.Abstractions/Services/Models/TargetEndpoint.cs +++ b/src/core/Elsa.Abstractions/Services/Models/TargetEndpoint.cs @@ -1,11 +1,7 @@ namespace Elsa.Services.Models { - public class TargetEndpoint : Endpoint + public class TargetEndpoint : Endpoint, ITargetEndpoint { - public TargetEndpoint() - { - } - public TargetEndpoint(IActivity activity) : base(activity) { } diff --git a/src/core/Elsa.Abstractions/Services/Models/Workflow.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs similarity index 56% rename from src/core/Elsa.Abstractions/Services/Models/Workflow.cs rename to src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs index ff0f49a75..594ffe5b9 100644 --- a/src/core/Elsa.Abstractions/Services/Models/Workflow.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs @@ -4,17 +4,17 @@ using Elsa.Models; namespace Elsa.Services.Models { - public class Workflow + public class WorkflowBlueprint : IWorkflowBlueprint { - public Workflow() + public WorkflowBlueprint() { - Activities = new List(); - Connections = new List(); - ActivityPropertyValueProviders = new Dictionary>(); + Activities = new List(); + Connections = new List(); + ActivityPropertyProviders = new ActivityPropertyProviders(); } - - public Workflow( - string workflowDefinitionId, + + public WorkflowBlueprint( + string definitionId, int version, bool isSingleton, bool isEnabled, @@ -24,11 +24,11 @@ namespace Elsa.Services.Models bool isPublished, WorkflowPersistenceBehavior persistenceBehavior, bool deleteCompletedInstances, - IEnumerable activities, - IEnumerable connections, - IDictionary> activityPropertyValueProviders) + IEnumerable activities, + IEnumerable connections, + IActivityPropertyProviders activityPropertyValueProviders) { - WorkflowDefinitionId = workflowDefinitionId; + DefinitionId = definitionId; Version = version; IsSingleton = isSingleton; IsEnabled = isEnabled; @@ -40,10 +40,11 @@ namespace Elsa.Services.Models DeleteCompletedInstances = deleteCompletedInstances; Activities = activities.ToList(); Connections = connections.ToList(); - ActivityPropertyValueProviders = activityPropertyValueProviders; + ActivityPropertyProviders = activityPropertyValueProviders; } - - public string WorkflowDefinitionId { get; set; } = default!; + + 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; } @@ -53,14 +54,10 @@ namespace Elsa.Services.Models public bool IsLatest { get; set; } public WorkflowPersistenceBehavior PersistenceBehavior { get; set; } public bool DeleteCompletedInstances { get; set; } - public ICollection Activities { get; set; } - - public ICollection Connections { get; set; } - - public IDictionary> ActivityPropertyValueProviders - { - get; - set; - } + public ICollection Activities { get; set; } + + public ICollection Connections { get; set; } + + public IActivityPropertyProviders ActivityPropertyProviders { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs index ced3496de..241c394cd 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs @@ -14,67 +14,60 @@ namespace Elsa.Services.Models { public WorkflowExecutionContext( IExpressionEvaluator expressionEvaluator, - IClock clock, IServiceProvider serviceProvider, - Workflow workflow, - WorkflowInstance workflowInstance, - WorkflowFault? workflowFault = default, - IEnumerable? executionLog = default) + WorkflowDefinition workflowDefinition, + WorkflowInstance workflowInstance + //IWorkflow workflow, + //WorkflowStatus status, + //Variables variables, + //string correlationId, + //IWorkflowFault? workflowFault, + //ICollection scheduledActivities, + //ICollection blockingActivities, + //IEnumerable? executionLog = default + ) { ServiceProvider = serviceProvider; - Workflow = workflow; + WorkflowDefinition = workflowDefinition; WorkflowInstance = workflowInstance; - CorrelationId = workflowInstance.CorrelationId; - Activities = workflow.Activities.ToList(); - Connections = workflow.Connections.ToList(); + //Workflow = workflow; + //CorrelationId = correlationId; ExpressionEvaluator = expressionEvaluator; - Clock = clock; - - var activityLookup = workflow.Activities.ToDictionary(x => x.Id); - - ScheduledActivities = new Stack( - workflowInstance.ScheduledActivities.Reverse().Select(x => CreateScheduledActivity(x, activityLookup))); - - BlockingActivities = new HashSet( - workflowInstance.BlockingActivities.Select(x => activityLookup[x.ActivityId])); - - Variables = workflowInstance.Variables; - Status = workflowInstance.Status; - PersistenceBehavior = workflow.PersistenceBehavior; - ActivityPropertyValueProviders = workflow.ActivityPropertyValueProviders; - WorkflowFault = workflowFault; - ExecutionLog = executionLog?.ToList() ?? new List(); + //ScheduledActivities = new Stack(scheduledActivities.Reverse()); + //BlockingActivities = new HashSet(blockingActivities); + //Variables = variables; + //Status = status; + //PersistenceBehavior = workflow.PersistenceBehavior; + //ActivityPropertyProviders = workflow.ActivityPropertyProviders; + //WorkflowFault = workflowFault; + //ExecutionLog = executionLog?.ToList() ?? new List(); IsFirstPass = true; } - - private ScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel, - IDictionary activityLookup) + + private IScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel, + IDictionary activityLookup) { var activity = activityLookup[scheduledActivityModel.ActivityId]; return new ScheduledActivity(activity, scheduledActivityModel.Input); } + public IWorkflowBlueprint WorkflowBlueprint { get; } public IServiceProvider ServiceProvider { get; } - public Workflow Workflow { get; } + public WorkflowDefinition WorkflowDefinition { get; } public WorkflowInstance WorkflowInstance { get; } - public string WorkflowDefinitionId => Workflow.WorkflowDefinitionId; - public string WorkflowInstanceId => WorkflowInstance.WorkflowInstanceId; - public int Version => WorkflowInstance.Version; - public ICollection Activities { get; } - public ICollection Connections { get; } public WorkflowStatus Status { get; set; } - public Stack ScheduledActivities { get; } - public HashSet BlockingActivities { get; } + public Stack ScheduledActivities { get; } + public HashSet BlockingActivities { get; } = new HashSet(new BlockingActivityEqualityComparer()); public Variables Variables { get; } public bool HasScheduledActivities => ScheduledActivities.Any(); - public ScheduledActivity? ScheduledActivity { get; private set; } - public WorkflowFault? WorkflowFault { get; private set; } + public IScheduledActivity? ScheduledActivity { get; private set; } + public IWorkflowFault? WorkflowFault { get; private set; } public object? Output { get; set; } - public void ScheduleActivities(IEnumerable activities, object? input = default) + public void ScheduleActivities(IEnumerable activityDefinitions, object? input = default) { - foreach (var activity in activities) - ScheduleActivity(activity, input); + foreach (var activityDefinition in activityDefinitions) + ScheduleActivity(activityDefinition, input); } public void ScheduleActivities(IEnumerable activities) @@ -83,28 +76,22 @@ namespace Elsa.Services.Models ScheduleActivity(activity); } - public void ScheduleActivity(IActivity activity, object? input = default) => - ScheduleActivity(new ScheduledActivity(activity, input)); + public void ScheduleActivity(ActivityDefinition activityDefinition, object? input = default) => + ScheduleActivity(new ScheduledActivity(activityDefinition, input)); public void ScheduleActivity(ScheduledActivity activity) => ScheduledActivities.Push(activity); - public ScheduledActivity PopScheduledActivity() => ScheduledActivity = ScheduledActivities.Pop(); - public ScheduledActivity PeekScheduledActivity() => ScheduledActivities.Peek(); + public IScheduledActivity PopScheduledActivity() => ScheduledActivity = ScheduledActivities.Pop(); + public IScheduledActivity PeekScheduledActivity() => ScheduledActivities.Peek(); public IExpressionEvaluator ExpressionEvaluator { get; } - public IClock Clock { get; } public string? CorrelationId { get; set; } - public WorkflowPersistenceBehavior PersistenceBehavior { get; set; } - - public IDictionary> ActivityPropertyValueProviders - { - get; - } - + public WorkflowPersistenceBehavior PersistenceBehavior { get; } + public IActivityPropertyProviders ActivityPropertyProviders { get; } public bool DeleteCompletedInstances { get; set; } - public ICollection ExecutionLog { get; } + public ICollection ExecutionLog { get; } public bool IsFirstPass { get; private set; } - public bool AddBlockingActivity(IActivity activity) => BlockingActivities.Add(activity); - public void SetVariable(string name, object? value) => Variables.Set(name, JToken.FromObject(value)); + public bool AddBlockingActivity(IActivity activity) => BlockingActivities.Add(new BlockingActivity(activity.Id, activity.Type)); + public void SetVariable(string name, object? value) => Variables.Set(name, JToken.FromObject(value!)); public T GetVariable(string name) => (T)GetVariable(name)!; public object? GetVariable(string name) => Variables.Get(name); public void CompletePass() => IsFirstPass = false; @@ -119,29 +106,26 @@ namespace Elsa.Services.Models public void Complete() => Status = WorkflowStatus.Completed; - public IActivity? GetActivity(string id) => Activities.FirstOrDefault(x => x.Id == id); + public IActivity? GetActivity(string id) => WorkflowBlueprint.Activities.FirstOrDefault(x => x.Id == id); - public WorkflowInstance UpdateWorkflowInstance() + public void UpdateWorkflowInstance(WorkflowInstance workflowInstance) { - var workflowInstance = WorkflowInstance; workflowInstance.Variables = Variables; workflowInstance.ScheduledActivities = new Stack( - ScheduledActivities.Select(x => new Elsa.Models.ScheduledActivity(x.Activity.Id, x.Input))); + ScheduledActivities.Select(x => new Elsa.Models.ScheduledActivity(x.ActivityDefinition.Id, x.Input))); - workflowInstance.Activities = Activities.Select(x => new ActivityInstance(x.Id, x.Type, x.Output, Serialize(x))).ToList(); - - workflowInstance.BlockingActivities = new HashSet( - BlockingActivities.Select(x => new BlockingActivity(x.Id, x.Type)), - new BlockingActivityEqualityComparer()); - + workflowInstance.Activities = + WorkflowBlueprint.Activities.Select(x => new ActivityInstance(x.Id, x.Type, x.Output, Serialize(x))).ToList(); + + workflowInstance.BlockingActivities = BlockingActivities; workflowInstance.Status = Status; workflowInstance.CorrelationId = CorrelationId; workflowInstance.Output = Output; var executionLog = workflowInstance.ExecutionLog.Concat( ExecutionLog.Select(x => new Elsa.Models.ExecutionLogEntry(x.Activity.Id, x.Timestamp))); - + workflowInstance.ExecutionLog = executionLog.ToList(); if (WorkflowFault != null) @@ -152,8 +136,6 @@ namespace Elsa.Services.Models Message = WorkflowFault.Message }; } - - return workflowInstance; } private JObject Serialize(IActivity activity) => JObject.FromObject(activity); diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowFault.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowFault.cs index 44eb888e0..9d6762541 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowFault.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowFault.cs @@ -2,7 +2,7 @@ namespace Elsa.Services.Models { - public class WorkflowFault + public class WorkflowFault : IWorkflowFault { public WorkflowFault(IActivity? activity = default, LocalizedString? message = default) { diff --git a/src/core/Elsa.Core/Builders/ActivityBuilder.cs b/src/core/Elsa.Core/Builders/ActivityBuilder.cs index 86a9ce9c7..979c1681a 100644 --- a/src/core/Elsa.Core/Builders/ActivityBuilder.cs +++ b/src/core/Elsa.Core/Builders/ActivityBuilder.cs @@ -1,5 +1,7 @@ using System; using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; using Elsa.Services; using Elsa.Services.Models; @@ -7,17 +9,26 @@ namespace Elsa.Builders { public class ActivityBuilder : IActivityBuilder { - public ActivityBuilder(IWorkflowBuilder workflowBuilder, - IActivity activity, + private readonly IActivityActivator _activityActivator; + + public ActivityBuilder( + Type activityType, + Action? setupActivity, + IWorkflowBuilder workflowBuilder, + IActivityActivator activityActivator, IDictionary? propertyValueProviders) { + _activityActivator = activityActivator; + ActivityType = activityType; + SetupActivity = setupActivity; WorkflowBuilder = workflowBuilder; - Activity = activity; PropertyValueProviders = propertyValueProviders; } + public Type ActivityType { get; } + public Action? SetupActivity { get; } public IWorkflowBuilder WorkflowBuilder { get; } - public IActivity Activity { get; } + public string? ActivityId { get; private set; } public IDictionary? PropertyValueProviders { get; } public IActivityBuilder Add( @@ -30,14 +41,14 @@ namespace Elsa.Builders Action>? setup = null, Action? branch = null) where T : class, IActivity => When(OutcomeNames.Done).Then(setup, branch); - + public IActivityBuilder Then( Action setup, Action? branch = null) where T : class, IActivity => When(OutcomeNames.Done).Then(setup, branch); - public IActivityBuilder Then(T activity, Action? branch = null) - where T : class, IActivity => When(OutcomeNames.Done).Then(activity, branch); + public IActivityBuilder Then(Action? branch = null) + where T : class, IActivity => When(OutcomeNames.Done).Then(branch); public IActivityBuilder Then(IActivityBuilder targetActivity) { @@ -45,7 +56,20 @@ namespace Elsa.Builders return this; } - public IActivity BuildActivity() => Activity; - public Workflow Build() => WorkflowBuilder.Build(); + public IActivityBuilder WithId(string id) + { + ActivityId = id; + return this; + } + + public Func> BuildActivityAsync() => + async (context, cancellationToken) => + { + var activity = _activityActivator.ActivateActivity(SetupActivity); + await context.SetActivityPropertiesAsync(activity, cancellationToken); + return activity; + }; + + public IWorkflowBlueprint Build() => WorkflowBuilder.Build(); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Builders/OutcomeBuilder.cs b/src/core/Elsa.Core/Builders/OutcomeBuilder.cs index cace469ea..231355b3e 100644 --- a/src/core/Elsa.Core/Builders/OutcomeBuilder.cs +++ b/src/core/Elsa.Core/Builders/OutcomeBuilder.cs @@ -37,6 +37,6 @@ namespace Elsa.Builders return activityBuilder; } - public Workflow Build() => WorkflowBuilder.Build(); + public WorkflowBlueprint Build() => WorkflowBuilder.Build(); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Builders/WorkflowBuilder.cs b/src/core/Elsa.Core/Builders/WorkflowBuilder.cs index a0f1ce538..5dfe8fa5a 100644 --- a/src/core/Elsa.Core/Builders/WorkflowBuilder.cs +++ b/src/core/Elsa.Core/Builders/WorkflowBuilder.cs @@ -10,17 +10,14 @@ namespace Elsa.Builders { public class WorkflowBuilder : IWorkflowBuilder { - private readonly IActivityResolver _activityResolver; private readonly IIdGenerator _idGenerator; private readonly IList _activityBuilders; private readonly IList _connectionBuilders; public WorkflowBuilder( - IActivityResolver activityResolver, IIdGenerator idGenerator, IServiceProvider serviceProvider) { - _activityResolver = activityResolver; _idGenerator = idGenerator; ServiceProvider = serviceProvider; Id = idGenerator.Generate(); @@ -86,65 +83,27 @@ namespace Elsa.Builders return this; } - public T BuildActivity(Action setup) where T : class, IActivity - { - var activity = _activityResolver.ResolveActivity(); - - setup(activity); - return activity; - } - - public Workflow Build() - { - var definitionId = !string.IsNullOrWhiteSpace(Id) ? Id : _idGenerator.Generate(); - var activities = _activityBuilders.Select(x => x.BuildActivity()).ToList(); - var connections = _connectionBuilders.Select(x => x.BuildConnection()).ToList(); - - // Generate deterministic activity ids. - var id = 1; - - foreach (var activity in activities.Where(activity => string.IsNullOrEmpty(activity.Id))) - activity.Id = $"activity-{id++}"; - - var activityPropertyValueProviders = _activityBuilders - .Select(x => (x.Activity.Id, x.PropertyValueProviders)) - .ToDictionary(x => x.Id, x => x.PropertyValueProviders!); - - var workflow = new Workflow( - definitionId, - Version, - IsSingleton, - false, - Name, - Description, - true, - true, - PersistenceBehavior, - DeleteCompletedInstances, - activities, - connections, - activityPropertyValueProviders); - - - return workflow; - } - - public IActivityBuilder New( - T activity, + public IActivityBuilder New( + Type activityType, + Action? setupActivity = default, Action? branch = default, IDictionary? propertyValueProviders = default) - where T : class, IActivity { - var activityBuilder = new ActivityBuilder(this, activity, propertyValueProviders); + var activityBuilder = new ActivityBuilder(activityType, setupActivity, this, propertyValueProviders); branch?.Invoke(activityBuilder); return activityBuilder; } + public IActivityBuilder New( + Action? branch = default, + IDictionary? propertyValueProviders = default) + where T : class, IActivity => + New(typeof(T), null, branch, propertyValueProviders); + public IActivityBuilder New( Action>? setup = default, Action? branch = default) where T : class, IActivity { - var activity = _activityResolver.ResolveActivity(); var propertyValuesBuilder = new SetupActivity(); setup?.Invoke(propertyValuesBuilder); @@ -152,18 +111,13 @@ namespace Elsa.Builders x => x.Key, x => (IActivityPropertyValueProvider)new DelegateActivityPropertyValueProvider(x.Value)); - return New(activity, branch, valueProviders); + return New(branch, valueProviders); } - + public IActivityBuilder New( - Action setup, - Action? branch = default) where T : class, IActivity - { - var activity = _activityResolver.ResolveActivity(); - setup(activity); - - return New(activity, branch); - } + Action? setup, + Action? branch = default) where T : class, IActivity => + New(typeof(T), x => setup?.Invoke((T)x), branch); public IActivityBuilder StartWith( Action>? setup = default, @@ -172,20 +126,18 @@ namespace Elsa.Builders var activityBuilder = New(setup, branch); return Add(activityBuilder, branch); } - + public IActivityBuilder StartWith( - Action setup, + Action? setup, Action? branch = default) where T : class, IActivity { var activityBuilder = New(setup, branch); return Add(activityBuilder, branch); } - public IActivityBuilder StartWith(T activity, Action? branch = default) - where T : class, IActivity - { - return Add(activity, branch); - } + public IActivityBuilder StartWith(Action? branch = default) + where T : class, IActivity => + Add(branch); public IActivityBuilder Add( Action>? setup = default, @@ -196,7 +148,7 @@ namespace Elsa.Builders } public IActivityBuilder Add( - Action setup, + Action setup, Action? branch = default) where T : class, IActivity { @@ -205,12 +157,11 @@ namespace Elsa.Builders } public IActivityBuilder Add( - T activity, Action? branch = default, IDictionary? propertyValueProviders = default) where T : class, IActivity { - var activityBuilder = new ActivityBuilder(this, activity, propertyValueProviders); + var activityBuilder = new ActivityBuilder(typeof(T), null, this, propertyValueProviders); return Add(activityBuilder); } @@ -249,22 +200,57 @@ namespace Elsa.Builders Action? branch = default) where T : class, IActivity => StartWith(setup, branch); - public IActivityBuilder Then(T activity, Action? branch = default) - where T : class, IActivity => StartWith(activity, branch); + public IActivityBuilder Then(Action? branch = default) + where T : class, IActivity => StartWith(branch); - public Workflow Build(IWorkflow workflow) + public IWorkflowBlueprint Build(IWorkflow workflow) { WithId(workflow.GetType().Name); workflow.Build(this); return Build(); } + + public IWorkflowBlueprint Build() + { + var definitionId = !string.IsNullOrWhiteSpace(Id) ? Id : _idGenerator.Generate(); + var activities = _activityBuilders.Select(x => new ActivityBlueprint(x.BuildActivityAsync())).ToList(); + var connections = _connectionBuilders.Select(x => x.BuildConnection()).ToList(); - public Workflow Build(Type workflowType) + // Generate deterministic activity ids. + var id = 1; + + foreach (var activity in activities.Where(activity => string.IsNullOrEmpty(activity.Id))) + activity.Id = $"activity-{id++}"; + + var activityPropertyValueProviders = _activityBuilders + .Select(x => (x.ActivityId, x.PropertyValueProviders)) + .ToDictionary(x => x.ActivityId!, x => x.PropertyValueProviders!); + + var workflow = new WorkflowBlueprint( + definitionId, + Version, + IsSingleton, + false, + Name, + Description, + true, + true, + PersistenceBehavior, + DeleteCompletedInstances, + activities, + connections, + new ActivityPropertyProviders(activityPropertyValueProviders)); + + + return workflow; + } + + public IWorkflowBlueprint Build(Type workflowType) { var workflow = (IWorkflow)ActivatorUtilities.GetServiceOrCreateInstance(ServiceProvider, workflowType); return Build(workflow); } - public Workflow Build() where T : IWorkflow => Build(typeof(T)); + public IWorkflowBlueprint Build() where T : IWorkflow => Build(typeof(T)); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Data/Services/DatabaseWorkflowProvider.cs b/src/core/Elsa.Core/Data/Services/DatabaseWorkflowProvider.cs index f15743116..8f449d08e 100644 --- a/src/core/Elsa.Core/Data/Services/DatabaseWorkflowProvider.cs +++ b/src/core/Elsa.Core/Data/Services/DatabaseWorkflowProvider.cs @@ -14,27 +14,27 @@ namespace Elsa.Data.Services public class DatabaseWorkflowProvider : IWorkflowProvider { private readonly IWorkflowDefinitionManager _workflowDefinitionManager; - private readonly IActivityResolver _activityResolver; + private readonly IActivityActivator _activityActivator; public DatabaseWorkflowProvider( IWorkflowDefinitionManager workflowDefinitionManager, - IActivityResolver activityResolver) + IActivityActivator activityActivator) { _workflowDefinitionManager = workflowDefinitionManager; - _activityResolver = activityResolver; + _activityActivator = activityActivator; } - public async Task> GetWorkflowsAsync(CancellationToken cancellationToken) + public async Task> GetWorkflowsAsync(CancellationToken cancellationToken) { var workflowDefinitions = await _workflowDefinitionManager.ListAsync(cancellationToken); return workflowDefinitions.Select(CreateWorkflow); } - private Workflow CreateWorkflow(WorkflowDefinition definition) + private WorkflowBlueprint CreateWorkflow(WorkflowDefinition definition) { var resolvedActivities = definition.Activities.Select(ResolveActivity).ToDictionary(x => x.Id); - var workflow = new Workflow( + var workflow = new WorkflowBlueprint( definition.WorkflowDefinitionVersionId, definition.Version, definition.IsSingleton, @@ -65,6 +65,6 @@ namespace Elsa.Data.Services } private IActivity ResolveActivity(ActivityDefinition activityDefinition) => - _activityResolver.ResolveActivity(activityDefinition); + _activityActivator.ActivateActivity(activityDefinition); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index 71df6f666..bb2b76899 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -88,6 +88,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton() .TryAddProvider(ServiceLifetime.Singleton) @@ -100,13 +101,13 @@ namespace Microsoft.Extensions.DependencyInjection .AddScoped() .AddSingleton() .AddScoped() - .AddSingleton() + .AddSingleton() .AddScoped() .AddScoped() .AddIndexProvider() .AddIndexProvider() .AddStartupRunner() - .AddTransient() + .AddTransient() .AddWorkflowProvider() .AddTransient() .AddTransient>(sp => sp.GetRequiredService) diff --git a/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs b/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs index ca68a9862..1e2dc1ab7 100644 --- a/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs +++ b/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs @@ -10,10 +10,10 @@ namespace Elsa.Extensions { public static class WorkflowRegistryExtensions { - public static Task GetWorkflowAsync(this IWorkflowRegistry workflowRegistry, CancellationToken cancellationToken) => + public static Task GetWorkflowAsync(this IWorkflowRegistry workflowRegistry, CancellationToken cancellationToken) => workflowRegistry.GetWorkflowAsync(typeof(T).Name, VersionOptions.Latest, cancellationToken); - public static async Task> GetWorkflowsByStartActivityAsync( + public static async Task> GetWorkflowsByStartActivityAsync( this IWorkflowRegistry workflowRegistry, CancellationToken cancellationToken = default) where T : IActivity @@ -22,7 +22,7 @@ namespace Elsa.Extensions return results.Select(x => (x.Workflow, (T)x.Activity)); } - public static async Task> GetWorkflowsByStartActivityAsync( + public static async Task> GetWorkflowsByStartActivityAsync( this IWorkflowRegistry workflowRegistry, string activityType, CancellationToken cancellationToken = default) diff --git a/src/core/Elsa.Core/Serialization/ActivitySerializer.cs b/src/core/Elsa.Core/Serialization/ActivitySerializer.cs new file mode 100644 index 000000000..b6aec1fac --- /dev/null +++ b/src/core/Elsa.Core/Serialization/ActivitySerializer.cs @@ -0,0 +1,13 @@ +using Elsa.Services; +using Newtonsoft.Json.Linq; + +namespace Elsa.Serialization +{ + public class ActivitySerializer : IActivitySerializer + { + private readonly ITokenSerializer _tokenSerializer; + public ActivitySerializer(ITokenSerializer tokenSerializer) => _tokenSerializer = tokenSerializer; + public JObject Serialize(T activity) where T : IActivity => _tokenSerializer.Serialize(activity); + public T Deserialize(JObject data) where T : IActivity => _tokenSerializer.Deserialize(data); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/ActivityResolver.cs b/src/core/Elsa.Core/Services/ActivityActivator.cs similarity index 64% rename from src/core/Elsa.Core/Services/ActivityResolver.cs rename to src/core/Elsa.Core/Services/ActivityActivator.cs index a3db14fab..2697fe7e7 100644 --- a/src/core/Elsa.Core/Services/ActivityResolver.cs +++ b/src/core/Elsa.Core/Services/ActivityActivator.cs @@ -1,60 +1,72 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using Microsoft.Extensions.DependencyInjection; - -namespace Elsa.Services -{ - public class ActivityResolver : IActivityResolver - { - private readonly IServiceProvider _serviceProvider; - private readonly Lazy> _lazyActivityTypeLookup; - - public ActivityResolver(IServiceProvider serviceProvider, Func> activitiesFunc) - { - _serviceProvider = serviceProvider; - _lazyActivityTypeLookup = new Lazy>( - () => - { - var activities = activitiesFunc(); - return activities.Select(x => x.GetType()).Distinct().ToDictionary(x => x.Name); - }); - } - - private IDictionary ActivityTypeLookup => _lazyActivityTypeLookup.Value; - - public Type ResolveActivityType(string activityTypeName) - { - if (!ActivityTypeLookup.ContainsKey(activityTypeName)) - { - var activityType = Type.GetType(activityTypeName); - - if (activityType == null) - throw new ArgumentException($"No such activity type: {activityTypeName}", nameof(activityTypeName)); - - ActivityTypeLookup[activityTypeName] = activityType; - } - - return ActivityTypeLookup[activityTypeName]; - } - - public IActivity ResolveActivity(string activityTypeName, Action? setup = null) - { - var activityType = ResolveActivityType(activityTypeName); - var activity = (IActivity)ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider, activityType); - - setup?.Invoke(activity); - return activity; - } - - public T ResolveActivity(Action? setup = null) where T : class, IActivity - { - var activity = ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider); - setup?.Invoke(activity); - return activity; - } - - public IEnumerable GetActivityTypes() => ActivityTypeLookup.Values.ToList(); - public Type? GetActivityType(string activityTypeName) => ActivityTypeLookup.ContainsKey(activityTypeName) ? ActivityTypeLookup[activityTypeName] : default; - } +using System; +using System.Collections.Generic; +using System.Linq; +using Elsa.Models; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Services +{ + public class ActivityActivator : IActivityActivator + { + private readonly IServiceProvider _serviceProvider; + private readonly Lazy> _lazyActivityTypeLookup; + + public ActivityActivator(IServiceProvider serviceProvider, Func> activitiesFunc) + { + _serviceProvider = serviceProvider; + _lazyActivityTypeLookup = new Lazy>( + () => + { + var activities = activitiesFunc(); + return activities.Select(x => x.GetType()).Distinct().ToDictionary(x => x.Name); + }); + } + + private IDictionary ActivityTypeLookup => _lazyActivityTypeLookup.Value; + + public Type GetActivityTypeByName(string activityTypeName) + { + if (!ActivityTypeLookup.ContainsKey(activityTypeName)) + { + var activityType = Type.GetType(activityTypeName); + + if (activityType == null) + throw new ArgumentException($"No such activity type: {activityTypeName}", nameof(activityTypeName)); + + ActivityTypeLookup[activityTypeName] = activityType; + } + + return ActivityTypeLookup[activityTypeName]; + } + + public IActivity ActivateActivity(string activityTypeName, Action? setup = null) + { + var activityType = GetActivityTypeByName(activityTypeName); + var activity = (IActivity)ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider, activityType); + + setup?.Invoke(activity); + return activity; + } + + public T ActivateActivity(Action? setup = null) where T : class, IActivity + { + var activity = ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider); + setup?.Invoke(activity); + return activity; + } + + public IActivity ActivateActivity(ActivityDefinition activityDefinition) + { + var activity = ActivateActivity(activityDefinition.Type); + activity.Description = activityDefinition.Description; + activity.Id = activityDefinition.Id; + activity.Name = activityDefinition.Name; + activity.DisplayName = activityDefinition.DisplayName; + activity.PersistWorkflow = activityDefinition.PersistWorkflow; + return activity; + } + + public IEnumerable GetActivityTypes() => ActivityTypeLookup.Values.ToList(); + public Type? GetActivityType(string activityTypeName) => ActivityTypeLookup.ContainsKey(activityTypeName) ? ActivityTypeLookup[activityTypeName] : default; + } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowActivator.cs b/src/core/Elsa.Core/Services/WorkflowActivator.cs deleted file mode 100644 index 70995c3b8..000000000 --- a/src/core/Elsa.Core/Services/WorkflowActivator.cs +++ /dev/null @@ -1,39 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.Models; -using Elsa.Services.Models; -using NodaTime; - -namespace Elsa.Services -{ - public class WorkflowActivator : IWorkflowActivator - { - private readonly IClock _clock; - private readonly IIdGenerator _idGenerator; - - public WorkflowActivator(IClock clock, IIdGenerator idGenerator) - { - _clock = clock; - _idGenerator = idGenerator; - } - - public Task ActivateAsync( - Workflow workflow, - string? correlationId = default, - CancellationToken cancellationToken = default) - { - var workflowInstance = new WorkflowInstance - { - WorkflowInstanceId = _idGenerator.Generate(), - WorkflowDefinitionId = workflow.WorkflowDefinitionId, - Version = workflow.Version, - Status = WorkflowStatus.Idle, - CorrelationId = correlationId, - CreatedAt = _clock.GetCurrentInstant(), - //Activities = workflow.Activities.Select(x => new ActivityInstance(x.Id, x.Type, x.Output, x.da)) - }; - - return Task.FromResult(workflowInstance); - } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowFactory.cs b/src/core/Elsa.Core/Services/WorkflowFactory.cs new file mode 100644 index 000000000..2067c2770 --- /dev/null +++ b/src/core/Elsa.Core/Services/WorkflowFactory.cs @@ -0,0 +1,47 @@ +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Models; +using Newtonsoft.Json.Linq; +using NodaTime; + +namespace Elsa.Services +{ + public class WorkflowFactory : IWorkflowFactory + { + private readonly IClock _clock; + private readonly IIdGenerator _idGenerator; + + public WorkflowFactory(IClock clock, IIdGenerator idGenerator) + { + _clock = clock; + _idGenerator = idGenerator; + } + + public Task InstantiateAsync( + WorkflowDefinition workflowDefinition, + string? correlationId = default, + CancellationToken cancellationToken = default) + { + var workflowInstance = new WorkflowInstance + { + WorkflowInstanceId = _idGenerator.Generate(), + WorkflowDefinitionId = workflowDefinition.WorkflowDefinitionId, + Version = workflowDefinition.Version, + Status = WorkflowStatus.Idle, + CorrelationId = correlationId, + CreatedAt = _clock.GetCurrentInstant(), + Activities = workflowDefinition.Activities.Select(CreateInstance).ToList(), + Variables = workflowDefinition.Variables != null ? new Variables(workflowDefinition.Variables) : new Variables(), + }; + + return Task.FromResult(workflowInstance); + } + + private ActivityInstance CreateInstance(ActivityDefinition activityDefinition) => new ActivityInstance( + activityDefinition.Id, + activityDefinition.Type, + null, + new JObject(activityDefinition.Data)); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowHost.cs b/src/core/Elsa.Core/Services/WorkflowHost.cs index 495b05d15..730d4b424 100644 --- a/src/core/Elsa.Core/Services/WorkflowHost.cs +++ b/src/core/Elsa.Core/Services/WorkflowHost.cs @@ -1,5 +1,4 @@ using System; -using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; @@ -11,15 +10,16 @@ using Elsa.Models; using Elsa.Queries; using Elsa.Services.Models; using MediatR; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using NodaTime; -using ScheduledActivity = Elsa.Services.Models.ScheduledActivity; +using YesSql.Data; namespace Elsa.Services { public class WorkflowHost : IWorkflowHost { - private delegate Task ActivityOperation( + private delegate ValueTask ActivityOperation( ActivityExecutionContext activityExecutionContext, IActivity activity, CancellationToken cancellationToken); @@ -30,110 +30,85 @@ namespace Elsa.Services private static readonly ActivityOperation Resume = (context, activity, cancellationToken) => activity.ResumeAsync(context, cancellationToken); + private readonly IWorkflowDefinitionManager _workflowDefinitionManager; private readonly IWorkflowInstanceManager _workflowInstanceManager; private readonly IWorkflowRegistry _workflowRegistry; - private readonly IWorkflowActivator _workflowActivator; + private readonly IWorkflowFactory _workflowFactory; private readonly IExpressionEvaluator _expressionEvaluator; - private readonly IClock _clock; private readonly IMediator _mediator; private readonly IServiceProvider _serviceProvider; private readonly ILogger _logger; public WorkflowHost( + IWorkflowDefinitionManager workflowDefinitionManager, IWorkflowInstanceManager workflowInstanceManager, IWorkflowRegistry workflowRegistry, - IWorkflowActivator workflowActivator, + IWorkflowFactory workflowFactory, IExpressionEvaluator expressionEvaluator, - IClock clock, IMediator mediator, IServiceProvider serviceProvider, ILogger logger) { + _workflowDefinitionManager = workflowDefinitionManager; _workflowInstanceManager = workflowInstanceManager; _workflowRegistry = workflowRegistry; - _workflowActivator = workflowActivator; + _workflowFactory = workflowFactory; _expressionEvaluator = expressionEvaluator; - _clock = clock; _mediator = mediator; _serviceProvider = serviceProvider; _logger = logger; } - public async Task RunWorkflowInstanceAsync( - string workflowInstanceId, + public async ValueTask RunWorkflowDefinitionAsync( + WorkflowDefinition workflowDefinition, string? activityId = default, object? input = default, + string? correlationId = default, CancellationToken cancellationToken = default) { - var workflowInstance = await _workflowInstanceManager.GetByWorkflowInstanceIdAsync(workflowInstanceId, cancellationToken); + var workflowInstance = await _workflowFactory.InstantiateAsync( + workflowDefinition, + correlationId, + cancellationToken); - if (workflowInstance == null) - { - _logger.LogDebug("Workflow instance {WorkflowInstanceId} does not exist.", workflowInstanceId); - return null; - } - - return await RunWorkflowInstanceAsync(workflowInstance, activityId, input, cancellationToken); + return await RunWorkflowInstanceAsync( + workflowDefinition, + workflowInstance, + activityId, + input, + cancellationToken); } - - public async Task RunWorkflowInstanceAsync( + + public async ValueTask RunWorkflowInstanceAsync( WorkflowInstance workflowInstance, string? activityId = default, object? input = default, CancellationToken cancellationToken = default) { - var workflow = await _workflowRegistry.GetWorkflowAsync( + var workflowDefinition = await _workflowDefinitionManager.GetAsync( workflowInstance.WorkflowDefinitionId, - VersionOptions.SpecificVersion(workflowInstance.Version), + VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken); + + if(workflowDefinition == null) + throw new WorkflowException($"No such workflow with definition {workflowInstance.WorkflowDefinitionId}"); + + return await RunWorkflowInstanceAsync( + workflowDefinition, + workflowInstance, + activityId, + input, cancellationToken); - - if (workflow == null) - throw new WorkflowException( - $"The specified workflow definition {workflowInstance.WorkflowDefinitionId} is either not registered or not published."); - - return await RunAsync(workflow, workflowInstance, activityId, input, cancellationToken); } - public async Task RunWorkflowDefinitionAsync( - string workflowDefinitionId, - string? activityId = default, - object? input = default, - string? correlationId = default, - CancellationToken cancellationToken = default) - { - var workflow = await _workflowRegistry.GetWorkflowAsync( - workflowDefinitionId, - VersionOptions.Published, - cancellationToken); - - if (workflow == null) - throw new WorkflowException( - $"The specified workflow definition {workflowDefinitionId} is either not registered or not published."); - - var workflowInstance = await _workflowActivator.ActivateAsync(workflow, correlationId, cancellationToken); - return await RunAsync(workflow, workflowInstance, activityId, input, cancellationToken); - } - - public async Task RunWorkflowAsync( - Workflow workflow, - string? activityId = default, - object? input = default, - string? correlationId = default, - CancellationToken cancellationToken = default) - { - var workflowInstance = await _workflowActivator.ActivateAsync(workflow, correlationId, cancellationToken); - return await RunAsync(workflow, workflowInstance, activityId, input, cancellationToken); - } - - private async Task RunAsync( - Workflow workflow, + public async ValueTask RunWorkflowInstanceAsync( + WorkflowDefinition workflowDefinition, WorkflowInstance workflowInstance, string? activityId = default, object? input = default, CancellationToken cancellationToken = default) { - var workflowExecutionContext = CreateWorkflowExecutionContext(workflow, workflowInstance); - var activity = activityId != null ? workflow.GetActivity(activityId) : default; + var workflowExecutionContext = CreateWorkflowExecutionContext(workflowDefinition, workflowInstance); + var activity = activityId != null ? workflowDefinition.GetActivityById(activityId) : default; switch (workflowExecutionContext.Status) { @@ -179,13 +154,14 @@ namespace Elsa.Services return workflowExecutionContext; } - private async Task BeginWorkflow(WorkflowExecutionContext workflowExecutionContext, - IActivity? activity, + private async Task BeginWorkflow( + WorkflowExecutionContext workflowExecutionContext, + ActivityDefinition? activity, object? input, CancellationToken cancellationToken) { if (activity == null) - activity = workflowExecutionContext.GetStartActivities().First(); + activity = workflowExecutionContext.WorkflowDefinition.GetStartActivities().First(); if (!await CanExecuteAsync(workflowExecutionContext, activity, input, cancellationToken)) return; @@ -201,30 +177,42 @@ namespace Elsa.Services await RunAsync(workflowExecutionContext, Execute, cancellationToken); } - private async Task ResumeWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, - IActivity activity, + private async Task ResumeWorkflowAsync( + WorkflowExecutionContext workflowExecutionContext, + ActivityDefinition activity, object? input, CancellationToken cancellationToken) { if (!await CanExecuteAsync(workflowExecutionContext, activity, input, cancellationToken)) return; - workflowExecutionContext.BlockingActivities.Remove(activity); + workflowExecutionContext.BlockingActivities.RemoveWhere(x => x.ActivityId == activity.Id); workflowExecutionContext.Status = WorkflowStatus.Running; workflowExecutionContext.ScheduleActivity(activity, input); await RunAsync(workflowExecutionContext, Resume, cancellationToken); } - private Task CanExecuteAsync(WorkflowExecutionContext workflowExecutionContext, - IActivity activity, + private async ValueTask CanExecuteAsync( + WorkflowExecutionContext workflowExecutionContext, + ActivityDefinition activityDefinition, object? input, CancellationToken cancellationToken) { - var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, activity, input); - return activity.CanExecuteAsync(activityExecutionContext, cancellationToken); + var activityExecutionContext = new ActivityExecutionContext( + workflowExecutionContext, + activityDefinition, + input); + + using var scope = _serviceProvider.CreateScope(); + var activityActivator = scope.ServiceProvider.GetRequiredService(); + var activity = await InstantiateActivityAsync( + activityActivator, + activityExecutionContext, + cancellationToken); + return await activity.CanExecuteAsync(activityExecutionContext, cancellationToken); } - private async Task RunAsync( + private async ValueTask RunAsync( WorkflowExecutionContext workflowExecutionContext, ActivityOperation activityOperation, CancellationToken cancellationToken = default) @@ -232,17 +220,27 @@ namespace Elsa.Services while (workflowExecutionContext.HasScheduledActivities) { var scheduledActivity = workflowExecutionContext.PopScheduledActivity(); - var currentActivity = scheduledActivity.Activity; + var currentActivity = scheduledActivity.ActivityDefinition; + var activityExecutionContext = new ActivityExecutionContext( workflowExecutionContext, currentActivity, scheduledActivity.Input); - await activityExecutionContext.SetActivityPropertiesAsync(cancellationToken); - var result = await activityOperation(activityExecutionContext, currentActivity, cancellationToken); - await _mediator.Publish(new ActivityExecuting(activityExecutionContext), cancellationToken); - await result.ExecuteAsync(activityExecutionContext, cancellationToken); - await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken); + using (var scope = _serviceProvider.CreateScope()) + { + var activityActivator = scope.ServiceProvider.GetRequiredService(); + + var activity = await InstantiateActivityAsync( + activityActivator, + 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(); @@ -252,34 +250,23 @@ namespace Elsa.Services workflowExecutionContext.Complete(); } - private ScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel, - IDictionary activityLookup) - { - var activity = activityLookup[scheduledActivityModel.ActivityId]; - return new ScheduledActivity(activity, scheduledActivityModel.Input); - } - private WorkflowExecutionContext CreateWorkflowExecutionContext( - Workflow workflow, - WorkflowInstance workflowInstance) - { - var activityInstanceLookup = workflowInstance.Activities.ToDictionary(x => x.Id); - - foreach (var activity in workflow.Activities) - { - if (!activityInstanceLookup.ContainsKey(activity.Id)) - continue; - - var activityInstance = activityInstanceLookup[activity.Id]; - activity.Output = activityInstance.Output; - } - - return new WorkflowExecutionContext( + WorkflowDefinition workflowDefinition, + WorkflowInstance workflowInstance) => + new WorkflowExecutionContext( _expressionEvaluator, - _clock, _serviceProvider, - workflow, + workflowDefinition, workflowInstance); + + private async ValueTask InstantiateActivityAsync( + IActivityActivator activityActivator, + ActivityExecutionContext activityExecutionContext, + CancellationToken cancellationToken) + { + var activity = activityActivator.ActivateActivity(activityExecutionContext.ActivityDefinition); + await activityExecutionContext.SetActivityPropertiesAsync(activity, cancellationToken); + return activity; } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowRegistry.cs b/src/core/Elsa.Core/Services/WorkflowRegistry.cs index 9c1460839..a44a81140 100644 --- a/src/core/Elsa.Core/Services/WorkflowRegistry.cs +++ b/src/core/Elsa.Core/Services/WorkflowRegistry.cs @@ -20,7 +20,7 @@ namespace Elsa.Services _serviceProvider = serviceProvider; } - public async Task> GetWorkflowsAsync(CancellationToken cancellationToken) + public async Task> GetWorkflowsAsync(CancellationToken cancellationToken) { using var scope = _serviceProvider.CreateScope(); var providers = scope.ServiceProvider.GetServices(); @@ -28,7 +28,7 @@ namespace Elsa.Services return tasks.SelectMany(x => x).ToList(); } - public async Task GetWorkflowAsync( + public async Task GetWorkflowAsync( string id, VersionOptions version, CancellationToken cancellationToken) @@ -36,7 +36,7 @@ namespace Elsa.Services var workflows = await GetWorkflowsAsync(cancellationToken).ToList(); return workflows - .Where(x => x.WorkflowDefinitionId == id) + .Where(x => x.Id == id) .OrderByDescending(x => x.Version) .WithVersion(version) .FirstOrDefault(); diff --git a/src/core/Elsa.Core/Services/WorkflowScheduler.cs b/src/core/Elsa.Core/Services/WorkflowScheduler.cs index ad2e7f2eb..5d6b4302e 100644 --- a/src/core/Elsa.Core/Services/WorkflowScheduler.cs +++ b/src/core/Elsa.Core/Services/WorkflowScheduler.cs @@ -22,20 +22,20 @@ namespace Elsa.Services { private readonly IBus _serviceBus; private readonly IWorkflowInstanceManager _workflowInstanceManager; - private readonly IWorkflowActivator _workflowActivator; + private readonly IWorkflowFactory _workflowFactory; private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowSchedulerQueue _queue; public WorkflowScheduler( IBus serviceBus, IWorkflowInstanceManager workflowInstanceManager, - IWorkflowActivator workflowActivator, + IWorkflowFactory workflowFactory, IWorkflowRegistry workflowRegistry, IWorkflowSchedulerQueue queue) { _serviceBus = serviceBus; _workflowInstanceManager = workflowInstanceManager; - _workflowActivator = workflowActivator; + _workflowFactory = workflowFactory; _workflowRegistry = workflowRegistry; _queue = queue; } @@ -103,7 +103,7 @@ namespace Elsa.Services where activity.Type == activityType select (workflow, activity); - var tuples = (IList<(Workflow Workflow, IActivity Activity)>)query.ToList(); + var tuples = (IList<(WorkflowBlueprint Workflow, IActivity Activity)>)query.ToList(); tuples = await FilterRunningSingletonsAsync(tuples).ToList(); @@ -118,7 +118,7 @@ namespace Elsa.Services } else { - var workflowInstance = await _workflowActivator.ActivateAsync( + var workflowInstance = await _workflowFactory.InstantiateAsync( workflow, correlationId, cancellationToken); @@ -164,19 +164,19 @@ namespace Elsa.Services cancellationToken); } - private async Task ScheduleWorkflowAsync(Workflow workflow, + private async Task ScheduleWorkflowAsync(WorkflowBlueprint workflowBlueprint, IActivity activity, object? input, string? correlationId, CancellationToken cancellationToken) { - var workflowInstance = await _workflowActivator.ActivateAsync(workflow, correlationId, cancellationToken); + var workflowInstance = await _workflowFactory.InstantiateAsync(workflowBlueprint, correlationId, cancellationToken); await _workflowInstanceManager.SaveAsync(workflowInstance, cancellationToken); await ScheduleWorkflowAsync(workflowInstance.WorkflowInstanceId, activity.Id, input, cancellationToken); } - private async Task> FilterRunningSingletonsAsync( - IEnumerable<(Workflow Workflow, IActivity Activity)> tuples) + private async Task> FilterRunningSingletonsAsync( + IEnumerable<(WorkflowBlueprint Workflow, IActivity Activity)> tuples) { var tupleList = tuples.ToList(); var transients = tupleList.Where(x => !x.Workflow.IsSingleton).ToList(); @@ -185,7 +185,7 @@ namespace Elsa.Services foreach (var tuple in singletons) { - var workflowDefinitionId = tuple.Workflow.WorkflowDefinitionId; + var workflowDefinitionId = tuple.Workflow.Id; var instances = await _workflowInstanceManager .ListByDefinitionAndStatusAsync(workflowDefinitionId, WorkflowStatus.Suspended); @@ -197,9 +197,9 @@ namespace Elsa.Services return result; } - private async Task> GetStartedWorkflowsAsync(Workflow workflow) + private async Task> GetStartedWorkflowsAsync(WorkflowBlueprint workflowBlueprint) { - var workflowDefinitionId = workflow.WorkflowDefinitionId; + var workflowDefinitionId = workflowBlueprint.Id; var suspendedInstances = await _workflowInstanceManager .ListByDefinitionAndStatusAsync(workflowDefinitionId, WorkflowStatus.Suspended); @@ -207,7 +207,7 @@ namespace Elsa.Services var idleInstances = await _workflowInstanceManager .ListByDefinitionAndStatusAsync(workflowDefinitionId, WorkflowStatus.Idle); - var startActivities = workflow.GetStartActivities().Select(x => x.Id).ToList(); + var startActivities = workflowBlueprint.GetStartActivities().Select(x => x.Id).ToList(); var startedInstances = suspendedInstances .Where(x => x.BlockingActivities.Any(y => startActivities.Contains(y.ActivityId))).ToList(); diff --git a/src/core/Elsa.Core/Services/WorkflowSchedulerQueue.cs b/src/core/Elsa.Core/Services/WorkflowSchedulerQueue.cs index 035ac2cc6..31376b4ed 100644 --- a/src/core/Elsa.Core/Services/WorkflowSchedulerQueue.cs +++ b/src/core/Elsa.Core/Services/WorkflowSchedulerQueue.cs @@ -5,15 +5,15 @@ namespace Elsa.Services { public class WorkflowSchedulerQueue : IWorkflowSchedulerQueue { - private readonly IDictionary<(string WorkflowDefinitionId, string ActivityId), (Workflow Workflow, IActivity Activity, object? Input, string? CorrelationId)> _nextWorkflowInstances; + private readonly IDictionary<(string WorkflowDefinitionId, string ActivityId), (WorkflowBlueprint Workflow, IActivity Activity, object? Input, string? CorrelationId)> _nextWorkflowInstances; public WorkflowSchedulerQueue() => - _nextWorkflowInstances = new Dictionary<(string WorkflowDefinitionId, string ActivityId), (Workflow Workflow, IActivity Activity, object? Input, string? CorrelationId)>(); + _nextWorkflowInstances = new Dictionary<(string WorkflowDefinitionId, string ActivityId), (WorkflowBlueprint Workflow, IActivity Activity, object? Input, string? CorrelationId)>(); - public void Enqueue(Workflow workflow, IActivity activity, object? input, string? correlationId) - => _nextWorkflowInstances[(workflow.WorkflowDefinitionId, activity.Id)] = (workflow, activity, input, correlationId); + public void Enqueue(WorkflowBlueprint workflowBlueprint, IActivity activity, object? input, string? correlationId) + => _nextWorkflowInstances[(workflowBlueprint.Id, activity.Id)] = (workflowBlueprint, activity, input, correlationId); - public (Workflow Workflow, IActivity Activity, object? Input, string? CorrelationId)? Dequeue(string workflowDefinitionId, string activityId) + public (WorkflowBlueprint Workflow, IActivity Activity, object? Input, string? CorrelationId)? Dequeue(string workflowDefinitionId, string activityId) { var key = (workflowDefinitionId, activityId); if(!_nextWorkflowInstances.ContainsKey(key)) diff --git a/src/core/Elsa.Core/WorkflowProviders/CodeWorkflowProvider.cs b/src/core/Elsa.Core/WorkflowProviders/CodeWorkflowProvider.cs index 0f8a50f34..489558325 100644 --- a/src/core/Elsa.Core/WorkflowProviders/CodeWorkflowProvider.cs +++ b/src/core/Elsa.Core/WorkflowProviders/CodeWorkflowProvider.cs @@ -23,7 +23,7 @@ namespace Elsa.WorkflowProviders _workflowBuilder = workflowBuilder; } - public Task> GetWorkflowsAsync(CancellationToken cancellationToken) => Task.FromResult(GetWorkflows()); - private IEnumerable GetWorkflows() => from workflow in _workflows let builder = _workflowBuilder() select builder.Build(workflow); + public Task> GetWorkflowsAsync(CancellationToken cancellationToken) => Task.FromResult(GetWorkflows()); + private IEnumerable GetWorkflows() => from workflow in _workflows let builder = _workflowBuilder() select builder.Build(workflow); } } \ No newline at end of file diff --git a/src/samples/Elsa.Samples.Serialization/Program.cs b/src/samples/Elsa.Samples.Serialization/Program.cs index 2b7c8e1af..0b494b150 100644 --- a/src/samples/Elsa.Samples.Serialization/Program.cs +++ b/src/samples/Elsa.Samples.Serialization/Program.cs @@ -33,7 +33,7 @@ namespace Elsa.Samples.Serialization Console.WriteLine(json); // Deserialize the workflow. - var deserializedWorkflow = serializer.Deserialize(json, SerializationFormats.Json); + var deserializedWorkflow = serializer.Deserialize(json, SerializationFormats.Json); // Get the workflow host. var workflowHost = services.GetService(); diff --git a/src/samples/Elsa.Samples.WorkflowDefinition/Program.cs b/src/samples/Elsa.Samples.WorkflowDefinition/Program.cs index 4e7d1e763..3b5011af0 100644 --- a/src/samples/Elsa.Samples.WorkflowDefinition/Program.cs +++ b/src/samples/Elsa.Samples.WorkflowDefinition/Program.cs @@ -22,12 +22,12 @@ namespace Elsa.Samples.WorkflowDefinition .AddSingleton(Console.In) .BuildServiceProvider(); - var activityResolver = services.GetRequiredService(); - var activity1 = activityResolver.ResolveActivity() + var activityResolver = services.GetRequiredService(); + var activity1 = activityResolver.ActivateActivity() .WithId("activity-1") .WithText("Hello world!"); - var activity2 = activityResolver.ResolveActivity() + var activity2 = activityResolver.ActivateActivity() .WithId("activity-2") .WithText("Goodbye cruel world...!"); diff --git a/src/server/Elsa.Server.GraphQL/Mapping/ActivityStateResolver.cs b/src/server/Elsa.Server.GraphQL/Mapping/ActivityStateResolver.cs index 84b015956..09cba324c 100644 --- a/src/server/Elsa.Server.GraphQL/Mapping/ActivityStateResolver.cs +++ b/src/server/Elsa.Server.GraphQL/Mapping/ActivityStateResolver.cs @@ -9,12 +9,12 @@ namespace Elsa.Server.GraphQL.Mapping public class ActivityStateResolver : IValueResolver { private readonly ITokenSerializer _serializer; - private readonly IActivityResolver _activityResolver; + private readonly IActivityActivator _activityActivator; - public ActivityStateResolver(ITokenSerializer serializer, IActivityResolver activityResolver) + public ActivityStateResolver(ITokenSerializer serializer, IActivityActivator activityActivator) { _serializer = serializer; - _activityResolver = activityResolver; + _activityActivator = activityActivator; } public Variables? Resolve(ActivityDefinitionInput source, ActivityDefinition destination, Variables? destMember, ResolutionContext context) diff --git a/src/server/Elsa.Server.GraphQL/Query.cs b/src/server/Elsa.Server.GraphQL/Query.cs index 077478556..819e38803 100644 --- a/src/server/Elsa.Server.GraphQL/Query.cs +++ b/src/server/Elsa.Server.GraphQL/Query.cs @@ -15,16 +15,16 @@ namespace Elsa.Server.GraphQL public class Query { public IEnumerable GetActivityDescriptors( - [Service] IActivityResolver activityResolver, + [Service] IActivityActivator activityActivator, [Service] IActivityDescriber describer) => - activityResolver.GetActivityTypes().Select(describer.Describe).ToList(); + activityActivator.GetActivityTypes().Select(describer.Describe).ToList(); public ActivityDescriptor? GetActivityDescriptor( - [Service] IActivityResolver activityResolver, + [Service] IActivityActivator activityActivator, [Service] IActivityDescriber describer, string typeName) { - var type = activityResolver.GetActivityType(typeName); + var type = activityActivator.GetActivityType(typeName); return type == null ? default : describer.Describe(type); } diff --git a/test/unit/Elsa.Core.UnitTests/WorkflowHostTests.cs b/test/unit/Elsa.Core.UnitTests/WorkflowHostTests.cs index 4b7cfb29a..8bd2782b9 100644 --- a/test/unit/Elsa.Core.UnitTests/WorkflowHostTests.cs +++ b/test/unit/Elsa.Core.UnitTests/WorkflowHostTests.cs @@ -32,7 +32,7 @@ namespace Elsa.Core.UnitTests _fixture = new Fixture().Customize(new NodaTimeCustomization()); _session = CreateSession(); - var workflowActivatorMock = new Mock(); + var workflowActivatorMock = new Mock(); var workflowRegistryMock = new Mock(); var workflowInstanceManager = new WorkflowInstanceManager(_session); var workflowExpressionEvaluatorMock = new Mock(); @@ -43,8 +43,8 @@ namespace Elsa.Core.UnitTests var serviceProvider = new ServiceCollection().BuildServiceProvider(); workflowActivatorMock - .Setup(x => x.ActivateAsync(It.IsAny(), It.IsAny(), It.IsAny())) - .ReturnsAsync((Workflow workflow, string? correlationId, CancellationToken cancellationToken) => new WorkflowInstance()); + .Setup(x => x.InstantiateAsync(It.IsAny(), It.IsAny(), It.IsAny())) + .ReturnsAsync((WorkflowBlueprint workflow, string? correlationId, CancellationToken cancellationToken) => new WorkflowInstance()); _workflowHost = new WorkflowHost( workflowInstanceManager, @@ -110,9 +110,9 @@ namespace Elsa.Core.UnitTests return activityMock.Object; } - private Workflow CreateWorkflow(IActivity activity) + private WorkflowBlueprint CreateWorkflow(IActivity activity) { - var workflow = new Workflow(); + var workflow = new WorkflowBlueprint(); workflow.Activities.Add(activity); return workflow; }