Incremental work on blueprints

This commit is contained in:
Sipke Schoorstra 2020-10-11 22:02:14 +02:00
parent 7a7d032e9e
commit d9d567a489
65 changed files with 821 additions and 622 deletions

View file

@ -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)
{

View file

@ -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);
}
}

View file

@ -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);
}

View file

@ -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();
}
}

View file

@ -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<string, IActivityPropertyValueProvider>? PropertyValueProviders { get; }
IActivityBuilder Add<T>(Action<ISetupActivity<T>>? setup = default) where T : class, IActivity;
IOutcomeBuilder When(string outcome);
IActivityBuilder Then(IActivityBuilder targetActivity);
IActivity BuildActivity();
Workflow Build();
IActivityBuilder WithId(string id);
Func<ActivityExecutionContext, CancellationToken, Task<IActivity>> BuildActivityAsync();
IWorkflowBlueprint Build();
}
}

View file

@ -7,6 +7,6 @@ namespace Elsa.Builders
{
IActivityBuilder Then<T>(Action<ISetupActivity<T>>? setup = default, Action<IActivityBuilder>? branch = default) where T : class, IActivity;
IActivityBuilder Then<T>(Action<T> setup, Action<IActivityBuilder>? branch = default) where T : class, IActivity;
IActivityBuilder Then<T>(T activity, Action<IActivityBuilder>? branch = default) where T : class, IActivity;
IActivityBuilder Then<T>(Action<IActivityBuilder>? branch = default) where T : class, IActivity;
}
}

View file

@ -7,6 +7,6 @@ namespace Elsa.Builders
IWorkflowBuilder WorkflowBuilder { get; }
IActivityBuilder Source { get; }
string? Outcome { get; }
Workflow Build();
WorkflowBlueprint Build();
}
}

View file

@ -24,7 +24,7 @@ namespace Elsa.Builders
IWorkflowBuilder WithDeleteCompletedInstances(bool value);
IWorkflowBuilder WithPersistenceBehavior(WorkflowPersistenceBehavior value);
IActivityBuilder New<T>(T activity,
IActivityBuilder New<T>(
Action<IActivityBuilder>? branch = default,
IDictionary<string, IActivityPropertyValueProvider>? propertyValueProviders = default)
where T : class, IActivity;
@ -36,7 +36,7 @@ namespace Elsa.Builders
Action<ISetupActivity<T>>? setup = default,
Action<IActivityBuilder>? branch = default) where T : class, IActivity;
IActivityBuilder StartWith<T>(T activity, Action<IActivityBuilder>? branch = default)
IActivityBuilder StartWith<T>(Action<IActivityBuilder>? branch = default)
where T : class, IActivity;
IActivityBuilder Add<T>(
@ -48,7 +48,6 @@ namespace Elsa.Builders
Action<IActivityBuilder>? branch = default) where T : class, IActivity;
IActivityBuilder Add<T>(
T activity,
Action<IActivityBuilder>? branch = default,
IDictionary<string, IActivityPropertyValueProvider>? propertyValueProviders = default)
where T : class, IActivity;
@ -62,9 +61,9 @@ namespace Elsa.Builders
Func<IActivityBuilder> target,
string outcome = OutcomeNames.Done);
Workflow Build();
Workflow Build(IWorkflow workflow);
Workflow Build(Type workflowType);
Workflow Build<T>() where T : IWorkflow;
IWorkflowBlueprint Build();
IWorkflowBlueprint Build(IWorkflow workflow);
IWorkflowBlueprint Build(Type workflowType);
IWorkflowBlueprint Build<T>() where T : IWorkflow;
}
}

View file

@ -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;
}
}
}

View file

@ -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<ActivityDefinition> 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;
}
}
}

View file

@ -7,17 +7,17 @@ namespace Elsa
{
public static class WorkflowExecutionContextExtensions
{
public static IEnumerable<IActivity> 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<IActivity> 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<Connection> GetInboundConnections(this WorkflowExecutionContext workflowExecutionContext, IActivity activity) => workflowExecutionContext.Connections.Where(x => x.Target.Activity == activity);
public static IEnumerable<Connection> GetOutboundConnections(this WorkflowExecutionContext workflowExecutionContext, IActivity activity) => workflowExecutionContext.Connections.Where(x => x.Source.Activity == activity);

View file

@ -8,10 +8,10 @@ namespace Elsa
{
public static class WorkflowExtensions
{
public static IEnumerable<Workflow> WithVersion(this IEnumerable<Workflow> query, VersionOptions version)
public static IEnumerable<WorkflowBlueprint> WithVersion(this IEnumerable<WorkflowBlueprint> query, VersionOptions version)
=> query.AsQueryable().WithVersion(version);
public static IQueryable<Workflow> WithVersion(this IQueryable<Workflow> query, VersionOptions version)
public static IQueryable<WorkflowBlueprint> WithVersion(this IQueryable<WorkflowBlueprint> 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<IActivity> GetStartActivities(this Workflow workflow)
public static IEnumerable<IActivity> 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<Connection> GetInboundConnections(this Workflow workflow, string activityId)
public static IEnumerable<Connection> 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<Connection> GetOutboundConnections(this Workflow workflow, string activityId)
public static IEnumerable<Connection> 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();
}
/// <summary>
/// Returns the full path of incoming activities.
/// </summary>
public static IEnumerable<string> GetInboundActivityPath(this Workflow workflow, string activityId)
public static IEnumerable<string> GetInboundActivityPath(this WorkflowBlueprint workflowBlueprint, string activityId)
{
var inspectedActivityIDs = new HashSet<string>();
return workflow.GetInboundActivityPathInternal(activityId, activityId, inspectedActivityIDs)
return workflowBlueprint.GetInboundActivityPathInternal(activityId, activityId, inspectedActivityIDs)
.Distinct().ToList();
}
private static IEnumerable<string> GetInboundActivityPathInternal(this Workflow workflowInstance,
private static IEnumerable<string> GetInboundActivityPathInternal(this WorkflowBlueprint workflowBlueprintBlueprintInstance,
string activityId,
string startingPointActivityId,
HashSet<string> 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);

View file

@ -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;
}
}

View file

@ -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<T> : ActivityDefinition where T : IActivity
{
public ActivityDefinition()
{
Type = typeof(T).Name;
}
public JObject Data { get; set; } = new JObject();
}
}

View file

@ -0,0 +1,11 @@
using Elsa.Services;
using Newtonsoft.Json.Linq;
namespace Elsa.Serialization
{
public interface IActivitySerializer
{
JObject Serialize<T>(T activity) where T : IActivity;
T Deserialize<T>(JObject data) where T : IActivity;
}
}

View file

@ -18,16 +18,16 @@ namespace Elsa.Services
public string? Description{ get; set; }
public bool PersistWorkflow { get; set; }
public Task<bool> CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(context, cancellationToken);
public Task<IActivityExecutionResult> ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnExecuteAsync(context, cancellationToken);
public ValueTask<bool> CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(context, cancellationToken);
public ValueTask<IActivityExecutionResult> ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnExecuteAsync(context, cancellationToken);
public Task<IActivityExecutionResult> ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnResumeAsync(context, cancellationToken);
public ValueTask<IActivityExecutionResult> ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnResumeAsync(context, cancellationToken);
protected virtual bool OnCanExecute(ActivityExecutionContext context) => OnCanExecute();
protected virtual bool OnCanExecute() => true;
protected virtual Task<bool> OnCanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(cancellationToken);
protected virtual Task<bool> OnCanExecuteAsync(CancellationToken cancellationToken) => Task.FromResult(OnCanExecute());
protected virtual Task<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => Task.FromResult(OnExecute(context));
protected virtual Task<IActivityExecutionResult> OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => Task.FromResult(OnResume(context));
protected virtual ValueTask<bool> OnCanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => OnCanExecuteAsync(cancellationToken);
protected virtual ValueTask<bool> OnCanExecuteAsync(CancellationToken cancellationToken) => new ValueTask<bool>(OnCanExecute());
protected virtual ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => new ValueTask<IActivityExecutionResult>(OnExecute(context));
protected virtual ValueTask<IActivityExecutionResult> OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) => new ValueTask<IActivityExecutionResult>(OnResume(context));
protected virtual IActivityExecutionResult OnExecute(ActivityExecutionContext context) => OnExecute();
protected virtual IActivityExecutionResult OnExecute() => Done();
protected virtual IActivityExecutionResult OnResume(ActivityExecutionContext context) => OnResume();

View file

@ -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<string, IDictionary<string, IActivityPropertyValueProvider>> _providers =
new Dictionary<string, IDictionary<string, IActivityPropertyValueProvider>>();
public ActivityPropertyProviders()
{
}
public ActivityPropertyProviders(IDictionary<string, IDictionary<string, IActivityPropertyValueProvider>> providers)
{
_providers = providers;
}
public void AddProvider(string activityId, string propertyName, IActivityPropertyValueProvider provider)
{
if (!_providers.TryGetValue(activityId, out var properties))
{
properties = new Dictionary<string, IActivityPropertyValueProvider>();
_providers.Add(activityId, properties);
}
properties[propertyName] = provider;
}
public IDictionary<string, IActivityPropertyValueProvider>? 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<ActivityPropertyAttribute>() != null;
}
}

View file

@ -13,7 +13,7 @@ namespace Elsa.Services
string Type { get; }
/// <summary>
/// Unique identifier of this activity.
/// Unique identifier of this activity within the workflow.
/// </summary>
string Id { get; set; }
@ -45,16 +45,16 @@ namespace Elsa.Services
/// <summary>
/// Returns a value of whether the specified activity can execute.
/// </summary>
Task<bool> CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default);
ValueTask<bool> CanExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default);
/// <summary>
/// Executes the specified activity.
/// </summary>
Task<IActivityExecutionResult> ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default);
ValueTask<IActivityExecutionResult> ExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default);
/// <summary>
/// Resumes the specified activity.
/// </summary>
Task<IActivityExecutionResult> ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default);
ValueTask<IActivityExecutionResult> ResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken = default);
}
}

View file

@ -0,0 +1,15 @@
using System;
using System.Collections.Generic;
using Elsa.Models;
namespace Elsa.Services
{
public interface IActivityActivator
{
IActivity ActivateActivity(string activityTypeName, Action<IActivity>? setup = default);
T ActivateActivity<T>(Action<T>? configure = default) where T : class, IActivity;
IActivity ActivateActivity(ActivityDefinition activityDefinition);
IEnumerable<Type> GetActivityTypes();
Type? GetActivityType(string activityTypeName);
}
}

View file

@ -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);
}
}

View file

@ -1,13 +0,0 @@
using System;
using System.Collections.Generic;
namespace Elsa.Services
{
public interface IActivityResolver
{
IActivity ResolveActivity(string activityTypeName, Action<IActivity>? setup = default);
T ResolveActivity<T>(Action<T>? configure = default) where T : class, IActivity;
IEnumerable<Type> GetActivityTypes();
Type? GetActivityType(string activityTypeName);
}
}

View file

@ -5,10 +5,10 @@ using Elsa.Services.Models;
namespace Elsa.Services
{
public interface IWorkflowActivator
public interface IWorkflowFactory
{
Task<WorkflowInstance> ActivateAsync(
Workflow workflow,
Task<WorkflowInstance> InstantiateAsync(
WorkflowDefinition workflowDefinition,
string? correlationId = default,
CancellationToken cancellationToken = default);
}

View file

@ -7,20 +7,24 @@ namespace Elsa.Services
{
public interface IWorkflowHost
{
Task<WorkflowExecutionContext> RunWorkflowAsync(Workflow workflow, string? activityId = default, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default);
ValueTask<WorkflowExecutionContext> RunWorkflowInstanceAsync(
WorkflowInstance workflowInstance,
string? activityId = default,
object? input = default,
CancellationToken cancellationToken = default);
ValueTask<WorkflowExecutionContext> RunWorkflowInstanceAsync(
WorkflowDefinition workflowDefinition,
WorkflowInstance workflowInstance,
string? activityId = default,
object? input = default,
CancellationToken cancellationToken = default);
Task<WorkflowExecutionContext?> RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default);
Task<WorkflowExecutionContext?> RunWorkflowInstanceAsync(WorkflowInstance workflowInstance, string? activityId = default, object? input = default, CancellationToken cancellationToken = default);
/// <summary>
/// Run a registered workflow by its ID.
/// </summary>
Task<WorkflowExecutionContext> RunWorkflowDefinitionAsync(string workflowDefinitionId, string? activityId = default, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default);
// /// <summary>
// /// Resume a workflow instance.
// /// </summary>
// Task ResumeAsync(string workflowInstanceId, string activityId, object? input = default, CancellationToken cancellationToken = default);
ValueTask<WorkflowExecutionContext> RunWorkflowDefinitionAsync(
WorkflowDefinition workflowDefinition,
string? activityId = default,
object? input = default,
string? correlationId = default,
CancellationToken cancellationToken = default);
}
}

View file

@ -10,6 +10,6 @@ namespace Elsa.Services
/// </summary>
public interface IWorkflowProvider
{
Task<IEnumerable<Workflow>> GetWorkflowsAsync(CancellationToken cancellationToken);
Task<IEnumerable<WorkflowBlueprint>> GetWorkflowsAsync(CancellationToken cancellationToken);
}
}

View file

@ -8,7 +8,7 @@ namespace Elsa.Services
{
public interface IWorkflowRegistry
{
Task<IEnumerable<Workflow>> GetWorkflowsAsync(CancellationToken cancellationToken = default);
Task<Workflow?> GetWorkflowAsync(string id, VersionOptions version, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowBlueprint>> GetWorkflowsAsync(CancellationToken cancellationToken = default);
Task<WorkflowBlueprint?> GetWorkflowAsync(string id, VersionOptions version, CancellationToken cancellationToken = default);
}
}

View file

@ -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);
}
}

View file

@ -0,0 +1,26 @@
using System;
using Elsa.ActivityResults;
namespace Elsa.Services.Models
{
public class ActivityBlueprint : IActivityBlueprint
{
public ActivityBlueprint()
{
}
public ActivityBlueprint(Func<ActivityExecutionContext, IActivity> createActivity)
{
CreateActivity = createActivity;
}
public ActivityBlueprint(string id, Func<ActivityExecutionContext, IActivity> createActivity)
{
Id = id;
CreateActivity = createActivity;
}
public string Id { get; set; } = default!;
public Func<ActivityExecutionContext, IActivity> CreateActivity { get; set; } = default!;
}
}

View file

@ -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<string>(0);
}
public WorkflowExecutionContext WorkflowExecutionContext { get; }
public IActivity Activity { get; }
public ActivityDefinition ActivityDefinition { get; }
public object? Input { get; }
public object? Output { get; set; }
public IReadOnlyCollection<string> Outcomes { get; set; }
@ -32,24 +33,11 @@ namespace Elsa.Services.Models
public T GetVariable<T>(string name) => WorkflowExecutionContext.GetVariable<T>(name);
public T GetService<T>() => WorkflowExecutionContext.ServiceProvider.GetService<T>();
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<ActivityPropertyAttribute>() != null;
public async ValueTask SetActivityPropertiesAsync(IActivity activity,
CancellationToken cancellationToken = default) =>
await WorkflowExecutionContext.ActivityPropertyProviders.SetActivityPropertiesAsync(
activity,
this,
cancellationToken);
}
}

View file

@ -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; }
}
}

View file

@ -2,7 +2,7 @@ using NodaTime;
namespace Elsa.Services.Models
{
public class ExecutionLogEntry
public class ExecutionLogEntry : IExecutionLogEntry
{
public ExecutionLogEntry(IActivity activity, Instant timestamp)
{

View file

@ -0,0 +1,10 @@
using System;
namespace Elsa.Services.Models
{
public interface IActivityBlueprint
{
public string Id { get; }
Func<ActivityExecutionContext, IActivity> CreateActivity { get; }
}
}

View file

@ -0,0 +1,8 @@
namespace Elsa.Services.Models
{
public interface IConnection
{
ISourceEndpoint Source { get; }
ITargetEndpoint Target { get; }
}
}

View file

@ -0,0 +1,7 @@
namespace Elsa.Services.Models
{
public interface IEndpoint
{
IActivity Activity { get; }
}
}

View file

@ -0,0 +1,10 @@
using NodaTime;
namespace Elsa.Services.Models
{
public interface IExecutionLogEntry
{
IActivity Activity { get; }
Instant Timestamp { get; }
}
}

View file

@ -0,0 +1,10 @@
using Elsa.Models;
namespace Elsa.Services.Models
{
public interface IScheduledActivity
{
ActivityDefinition ActivityDefinition { get; }
object? Input { get; }
}
}

View file

@ -0,0 +1,7 @@
namespace Elsa.Services.Models
{
public interface ISourceEndpoint : IEndpoint
{
string Outcome { get; }
}
}

View file

@ -0,0 +1,6 @@
namespace Elsa.Services.Models
{
public interface ITargetEndpoint : IEndpoint
{
}
}

View file

@ -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<IActivityBlueprint> Activities { get; }
public ICollection<IConnection> Connections { get; }
IActivityPropertyProviders ActivityPropertyProviders { get; }
}
}

View file

@ -0,0 +1,10 @@
using Microsoft.Extensions.Localization;
namespace Elsa.Services.Models
{
public interface IWorkflowFault
{
IActivity? FaultedActivity { get; }
LocalizedString? Message { get; }
}
}

View file

@ -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; }
}
}

View file

@ -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; }
}
}

View file

@ -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)
{
}

View file

@ -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<IActivity>();
Connections = new List<Connection>();
ActivityPropertyValueProviders = new Dictionary<string, IDictionary<string, IActivityPropertyValueProvider>>();
Activities = new List<IActivityBlueprint>();
Connections = new List<IConnection>();
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<IActivity> activities,
IEnumerable<Connection> connections,
IDictionary<string, IDictionary<string, IActivityPropertyValueProvider>> activityPropertyValueProviders)
IEnumerable<IActivityBlueprint> activities,
IEnumerable<IConnection> 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<IActivity> Activities { get; set; }
public ICollection<Connection> Connections { get; set; }
public IDictionary<string, IDictionary<string, IActivityPropertyValueProvider>> ActivityPropertyValueProviders
{
get;
set;
}
public ICollection<IActivityBlueprint> Activities { get; set; }
public ICollection<IConnection> Connections { get; set; }
public IActivityPropertyProviders ActivityPropertyProviders { get; set; }
}
}

View file

@ -14,67 +14,60 @@ namespace Elsa.Services.Models
{
public WorkflowExecutionContext(
IExpressionEvaluator expressionEvaluator,
IClock clock,
IServiceProvider serviceProvider,
Workflow workflow,
WorkflowInstance workflowInstance,
WorkflowFault? workflowFault = default,
IEnumerable<ExecutionLogEntry>? executionLog = default)
WorkflowDefinition workflowDefinition,
WorkflowInstance workflowInstance
//IWorkflow workflow,
//WorkflowStatus status,
//Variables variables,
//string correlationId,
//IWorkflowFault? workflowFault,
//ICollection<IScheduledActivity> scheduledActivities,
//ICollection<IActivity> blockingActivities,
//IEnumerable<IExecutionLogEntry>? 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<ScheduledActivity>(
workflowInstance.ScheduledActivities.Reverse().Select(x => CreateScheduledActivity(x, activityLookup)));
BlockingActivities = new HashSet<IActivity>(
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<ExecutionLogEntry>();
//ScheduledActivities = new Stack<IScheduledActivity>(scheduledActivities.Reverse());
//BlockingActivities = new HashSet<IActivity>(blockingActivities);
//Variables = variables;
//Status = status;
//PersistenceBehavior = workflow.PersistenceBehavior;
//ActivityPropertyProviders = workflow.ActivityPropertyProviders;
//WorkflowFault = workflowFault;
//ExecutionLog = executionLog?.ToList() ?? new List<IExecutionLogEntry>();
IsFirstPass = true;
}
private ScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel,
IDictionary<string, IActivity> activityLookup)
private IScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel,
IDictionary<string, ActivityDefinition> 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<IActivity> Activities { get; }
public ICollection<Connection> Connections { get; }
public WorkflowStatus Status { get; set; }
public Stack<ScheduledActivity> ScheduledActivities { get; }
public HashSet<IActivity> BlockingActivities { get; }
public Stack<IScheduledActivity> ScheduledActivities { get; }
public HashSet<BlockingActivity> BlockingActivities { get; } = new HashSet<BlockingActivity>(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<IActivity> activities, object? input = default)
public void ScheduleActivities(IEnumerable<ActivityDefinition> activityDefinitions, object? input = default)
{
foreach (var activity in activities)
ScheduleActivity(activity, input);
foreach (var activityDefinition in activityDefinitions)
ScheduleActivity(activityDefinition, input);
}
public void ScheduleActivities(IEnumerable<ScheduledActivity> 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<string, IDictionary<string, IActivityPropertyValueProvider>> ActivityPropertyValueProviders
{
get;
}
public WorkflowPersistenceBehavior PersistenceBehavior { get; }
public IActivityPropertyProviders ActivityPropertyProviders { get; }
public bool DeleteCompletedInstances { get; set; }
public ICollection<ExecutionLogEntry> ExecutionLog { get; }
public ICollection<IExecutionLogEntry> 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<T>(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<Elsa.Models.ScheduledActivity>(
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<BlockingActivity>(
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);

View file

@ -2,7 +2,7 @@
namespace Elsa.Services.Models
{
public class WorkflowFault
public class WorkflowFault : IWorkflowFault
{
public WorkflowFault(IActivity? activity = default, LocalizedString? message = default)
{

View file

@ -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<IActivity>? setupActivity,
IWorkflowBuilder workflowBuilder,
IActivityActivator activityActivator,
IDictionary<string, IActivityPropertyValueProvider>? propertyValueProviders)
{
_activityActivator = activityActivator;
ActivityType = activityType;
SetupActivity = setupActivity;
WorkflowBuilder = workflowBuilder;
Activity = activity;
PropertyValueProviders = propertyValueProviders;
}
public Type ActivityType { get; }
public Action<IActivity>? SetupActivity { get; }
public IWorkflowBuilder WorkflowBuilder { get; }
public IActivity Activity { get; }
public string? ActivityId { get; private set; }
public IDictionary<string, IActivityPropertyValueProvider>? PropertyValueProviders { get; }
public IActivityBuilder Add<T>(
@ -30,14 +41,14 @@ namespace Elsa.Builders
Action<ISetupActivity<T>>? setup = null,
Action<IActivityBuilder>? branch = null)
where T : class, IActivity => When(OutcomeNames.Done).Then(setup, branch);
public IActivityBuilder Then<T>(
Action<T> setup,
Action<IActivityBuilder>? branch = null)
where T : class, IActivity => When(OutcomeNames.Done).Then(setup, branch);
public IActivityBuilder Then<T>(T activity, Action<IActivityBuilder>? branch = null)
where T : class, IActivity => When(OutcomeNames.Done).Then(activity, branch);
public IActivityBuilder Then<T>(Action<IActivityBuilder>? branch = null)
where T : class, IActivity => When(OutcomeNames.Done).Then<T>(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<ActivityExecutionContext, CancellationToken, Task<IActivity>> BuildActivityAsync() =>
async (context, cancellationToken) =>
{
var activity = _activityActivator.ActivateActivity(SetupActivity);
await context.SetActivityPropertiesAsync(activity, cancellationToken);
return activity;
};
public IWorkflowBlueprint Build() => WorkflowBuilder.Build();
}
}

View file

@ -37,6 +37,6 @@ namespace Elsa.Builders
return activityBuilder;
}
public Workflow Build() => WorkflowBuilder.Build();
public WorkflowBlueprint Build() => WorkflowBuilder.Build();
}
}

View file

@ -10,17 +10,14 @@ namespace Elsa.Builders
{
public class WorkflowBuilder : IWorkflowBuilder
{
private readonly IActivityResolver _activityResolver;
private readonly IIdGenerator _idGenerator;
private readonly IList<IActivityBuilder> _activityBuilders;
private readonly IList<IConnectionBuilder> _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<T>(Action<T> setup) where T : class, IActivity
{
var activity = _activityResolver.ResolveActivity<T>();
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>(
T activity,
public IActivityBuilder New(
Type activityType,
Action<IActivity>? setupActivity = default,
Action<IActivityBuilder>? branch = default,
IDictionary<string, IActivityPropertyValueProvider>? 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<T>(
Action<IActivityBuilder>? branch = default,
IDictionary<string, IActivityPropertyValueProvider>? propertyValueProviders = default)
where T : class, IActivity =>
New(typeof(T), null, branch, propertyValueProviders);
public IActivityBuilder New<T>(
Action<ISetupActivity<T>>? setup = default,
Action<IActivityBuilder>? branch = default) where T : class, IActivity
{
var activity = _activityResolver.ResolveActivity<T>();
var propertyValuesBuilder = new SetupActivity<T>();
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<T>(branch, valueProviders);
}
public IActivityBuilder New<T>(
Action<T> setup,
Action<IActivityBuilder>? branch = default) where T : class, IActivity
{
var activity = _activityResolver.ResolveActivity<T>();
setup(activity);
return New(activity, branch);
}
Action<T>? setup,
Action<IActivityBuilder>? branch = default) where T : class, IActivity =>
New(typeof(T), x => setup?.Invoke((T)x), branch);
public IActivityBuilder StartWith<T>(
Action<ISetupActivity<T>>? setup = default,
@ -172,20 +126,18 @@ namespace Elsa.Builders
var activityBuilder = New(setup, branch);
return Add(activityBuilder, branch);
}
public IActivityBuilder StartWith<T>(
Action<T> setup,
Action<T>? setup,
Action<IActivityBuilder>? branch = default) where T : class, IActivity
{
var activityBuilder = New(setup, branch);
return Add(activityBuilder, branch);
}
public IActivityBuilder StartWith<T>(T activity, Action<IActivityBuilder>? branch = default)
where T : class, IActivity
{
return Add(activity, branch);
}
public IActivityBuilder StartWith<T>(Action<IActivityBuilder>? branch = default)
where T : class, IActivity =>
Add<T>(branch);
public IActivityBuilder Add<T>(
Action<ISetupActivity<T>>? setup = default,
@ -196,7 +148,7 @@ namespace Elsa.Builders
}
public IActivityBuilder Add<T>(
Action<T> setup,
Action<T> setup,
Action<IActivityBuilder>? branch = default)
where T : class, IActivity
{
@ -205,12 +157,11 @@ namespace Elsa.Builders
}
public IActivityBuilder Add<T>(
T activity,
Action<IActivityBuilder>? branch = default,
IDictionary<string, IActivityPropertyValueProvider>? 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<IActivityBuilder>? branch = default)
where T : class, IActivity => StartWith(setup, branch);
public IActivityBuilder Then<T>(T activity, Action<IActivityBuilder>? branch = default)
where T : class, IActivity => StartWith(activity, branch);
public IActivityBuilder Then<T>(Action<IActivityBuilder>? branch = default)
where T : class, IActivity => StartWith<T>(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<T>() where T : IWorkflow => Build(typeof(T));
public IWorkflowBlueprint Build<T>() where T : IWorkflow => Build(typeof(T));
}
}

View file

@ -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<IEnumerable<Workflow>> GetWorkflowsAsync(CancellationToken cancellationToken)
public async Task<IEnumerable<WorkflowBlueprint>> 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);
}
}

View file

@ -88,6 +88,7 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton<ITokenSerializerProvider, TokenSerializerProvider>()
.AddSingleton<ITokenSerializer, TokenSerializer>()
.AddSingleton<IWorkflowSerializer, WorkflowSerializer>()
.AddSingleton<IActivitySerializer, ActivitySerializer>()
.AddSingleton<TypeConverter>()
.AddSingleton<ITypeMap, TypeMap>()
.TryAddProvider<ITokenFormatter, JsonTokenFormatter>(ServiceLifetime.Singleton)
@ -100,13 +101,13 @@ namespace Microsoft.Extensions.DependencyInjection
.AddScoped<IWorkflowScheduler, WorkflowScheduler>()
.AddSingleton<IWorkflowSchedulerQueue, WorkflowSchedulerQueue>()
.AddScoped<IWorkflowHost, WorkflowHost>()
.AddSingleton<IWorkflowActivator, WorkflowActivator>()
.AddSingleton<IWorkflowFactory, WorkflowFactory>()
.AddScoped<IWorkflowDefinitionManager, WorkflowDefinitionManager>()
.AddScoped<IWorkflowInstanceManager, WorkflowInstanceManager>()
.AddIndexProvider<WorkflowDefinitionIndexProvider>()
.AddIndexProvider<WorkflowInstanceIndexProvider>()
.AddStartupRunner()
.AddTransient<IActivityResolver, ActivityResolver>()
.AddTransient<IActivityActivator, ActivityActivator>()
.AddWorkflowProvider<CodeWorkflowProvider>()
.AddTransient<IWorkflowBuilder, WorkflowBuilder>()
.AddTransient<Func<IWorkflowBuilder>>(sp => sp.GetRequiredService<IWorkflowBuilder>)

View file

@ -10,10 +10,10 @@ namespace Elsa.Extensions
{
public static class WorkflowRegistryExtensions
{
public static Task<Workflow> GetWorkflowAsync<T>(this IWorkflowRegistry workflowRegistry, CancellationToken cancellationToken) =>
public static Task<WorkflowBlueprint> GetWorkflowAsync<T>(this IWorkflowRegistry workflowRegistry, CancellationToken cancellationToken) =>
workflowRegistry.GetWorkflowAsync(typeof(T).Name, VersionOptions.Latest, cancellationToken);
public static async Task<IEnumerable<(Workflow Workflow, T Activity)>> GetWorkflowsByStartActivityAsync<T>(
public static async Task<IEnumerable<(WorkflowBlueprint Workflow, T Activity)>> GetWorkflowsByStartActivityAsync<T>(
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<IEnumerable<(Workflow Workflow, IActivity Activity)>> GetWorkflowsByStartActivityAsync(
public static async Task<IEnumerable<(WorkflowBlueprint Workflow, IActivity Activity)>> GetWorkflowsByStartActivityAsync(
this IWorkflowRegistry workflowRegistry,
string activityType,
CancellationToken cancellationToken = default)

View file

@ -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>(T activity) where T : IActivity => _tokenSerializer.Serialize(activity);
public T Deserialize<T>(JObject data) where T : IActivity => _tokenSerializer.Deserialize<T>(data);
}
}

View file

@ -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<IDictionary<string, Type>> _lazyActivityTypeLookup;
public ActivityResolver(IServiceProvider serviceProvider, Func<IEnumerable<IActivity>> activitiesFunc)
{
_serviceProvider = serviceProvider;
_lazyActivityTypeLookup = new Lazy<IDictionary<string, Type>>(
() =>
{
var activities = activitiesFunc();
return activities.Select(x => x.GetType()).Distinct().ToDictionary(x => x.Name);
});
}
private IDictionary<string, Type> 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<IActivity>? setup = null)
{
var activityType = ResolveActivityType(activityTypeName);
var activity = (IActivity)ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider, activityType);
setup?.Invoke(activity);
return activity;
}
public T ResolveActivity<T>(Action<T>? setup = null) where T : class, IActivity
{
var activity = ActivatorUtilities.GetServiceOrCreateInstance<T>(_serviceProvider);
setup?.Invoke(activity);
return activity;
}
public IEnumerable<Type> 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<IDictionary<string, Type>> _lazyActivityTypeLookup;
public ActivityActivator(IServiceProvider serviceProvider, Func<IEnumerable<IActivity>> activitiesFunc)
{
_serviceProvider = serviceProvider;
_lazyActivityTypeLookup = new Lazy<IDictionary<string, Type>>(
() =>
{
var activities = activitiesFunc();
return activities.Select(x => x.GetType()).Distinct().ToDictionary(x => x.Name);
});
}
private IDictionary<string, Type> 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<IActivity>? setup = null)
{
var activityType = GetActivityTypeByName(activityTypeName);
var activity = (IActivity)ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider, activityType);
setup?.Invoke(activity);
return activity;
}
public T ActivateActivity<T>(Action<T>? setup = null) where T : class, IActivity
{
var activity = ActivatorUtilities.GetServiceOrCreateInstance<T>(_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<Type> GetActivityTypes() => ActivityTypeLookup.Values.ToList();
public Type? GetActivityType(string activityTypeName) => ActivityTypeLookup.ContainsKey(activityTypeName) ? ActivityTypeLookup[activityTypeName] : default;
}
}

View file

@ -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<WorkflowInstance> 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);
}
}
}

View file

@ -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<WorkflowInstance> 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));
}
}

View file

@ -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<IActivityExecutionResult> ActivityOperation(
private delegate ValueTask<IActivityExecutionResult> 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<WorkflowHost> logger)
{
_workflowDefinitionManager = workflowDefinitionManager;
_workflowInstanceManager = workflowInstanceManager;
_workflowRegistry = workflowRegistry;
_workflowActivator = workflowActivator;
_workflowFactory = workflowFactory;
_expressionEvaluator = expressionEvaluator;
_clock = clock;
_mediator = mediator;
_serviceProvider = serviceProvider;
_logger = logger;
}
public async Task<WorkflowExecutionContext?> RunWorkflowInstanceAsync(
string workflowInstanceId,
public async ValueTask<WorkflowExecutionContext> 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<WorkflowExecutionContext?> RunWorkflowInstanceAsync(
public async ValueTask<WorkflowExecutionContext> 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<WorkflowExecutionContext> 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<WorkflowExecutionContext> 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<WorkflowExecutionContext> RunAsync(
Workflow workflow,
public async ValueTask<WorkflowExecutionContext> 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<bool> CanExecuteAsync(WorkflowExecutionContext workflowExecutionContext,
IActivity activity,
private async ValueTask<bool> 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<IActivityActivator>();
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<IActivityActivator>();
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<string, IActivity> 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<IActivity> InstantiateActivityAsync(
IActivityActivator activityActivator,
ActivityExecutionContext activityExecutionContext,
CancellationToken cancellationToken)
{
var activity = activityActivator.ActivateActivity(activityExecutionContext.ActivityDefinition);
await activityExecutionContext.SetActivityPropertiesAsync(activity, cancellationToken);
return activity;
}
}
}

View file

@ -20,7 +20,7 @@ namespace Elsa.Services
_serviceProvider = serviceProvider;
}
public async Task<IEnumerable<Workflow>> GetWorkflowsAsync(CancellationToken cancellationToken)
public async Task<IEnumerable<WorkflowBlueprint>> GetWorkflowsAsync(CancellationToken cancellationToken)
{
using var scope = _serviceProvider.CreateScope();
var providers = scope.ServiceProvider.GetServices<IWorkflowProvider>();
@ -28,7 +28,7 @@ namespace Elsa.Services
return tasks.SelectMany(x => x).ToList();
}
public async Task<Workflow?> GetWorkflowAsync(
public async Task<WorkflowBlueprint?> 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();

View file

@ -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<IEnumerable<(Workflow, IActivity)>> FilterRunningSingletonsAsync(
IEnumerable<(Workflow Workflow, IActivity Activity)> tuples)
private async Task<IEnumerable<(WorkflowBlueprint, IActivity)>> 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<IEnumerable<WorkflowInstance>> GetStartedWorkflowsAsync(Workflow workflow)
private async Task<IEnumerable<WorkflowInstance>> 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();

View file

@ -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))

View file

@ -23,7 +23,7 @@ namespace Elsa.WorkflowProviders
_workflowBuilder = workflowBuilder;
}
public Task<IEnumerable<Workflow>> GetWorkflowsAsync(CancellationToken cancellationToken) => Task.FromResult(GetWorkflows());
private IEnumerable<Workflow> GetWorkflows() => from workflow in _workflows let builder = _workflowBuilder() select builder.Build(workflow);
public Task<IEnumerable<WorkflowBlueprint>> GetWorkflowsAsync(CancellationToken cancellationToken) => Task.FromResult(GetWorkflows());
private IEnumerable<WorkflowBlueprint> GetWorkflows() => from workflow in _workflows let builder = _workflowBuilder() select builder.Build(workflow);
}
}

View file

@ -33,7 +33,7 @@ namespace Elsa.Samples.Serialization
Console.WriteLine(json);
// Deserialize the workflow.
var deserializedWorkflow = serializer.Deserialize<Workflow>(json, SerializationFormats.Json);
var deserializedWorkflow = serializer.Deserialize<WorkflowBlueprint>(json, SerializationFormats.Json);
// Get the workflow host.
var workflowHost = services.GetService<IWorkflowHost>();

View file

@ -22,12 +22,12 @@ namespace Elsa.Samples.WorkflowDefinition
.AddSingleton(Console.In)
.BuildServiceProvider();
var activityResolver = services.GetRequiredService<IActivityResolver>();
var activity1 = activityResolver.ResolveActivity<WriteLine>()
var activityResolver = services.GetRequiredService<IActivityActivator>();
var activity1 = activityResolver.ActivateActivity<WriteLine>()
.WithId("activity-1")
.WithText("Hello world!");
var activity2 = activityResolver.ResolveActivity<WriteLine>()
var activity2 = activityResolver.ActivateActivity<WriteLine>()
.WithId("activity-2")
.WithText("Goodbye cruel world...!");

View file

@ -9,12 +9,12 @@ namespace Elsa.Server.GraphQL.Mapping
public class ActivityStateResolver : IValueResolver<ActivityDefinitionInput, ActivityDefinition, Variables?>
{
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)

View file

@ -15,16 +15,16 @@ namespace Elsa.Server.GraphQL
public class Query
{
public IEnumerable<ActivityDescriptor> 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);
}

View file

@ -32,7 +32,7 @@ namespace Elsa.Core.UnitTests
_fixture = new Fixture().Customize(new NodaTimeCustomization());
_session = CreateSession();
var workflowActivatorMock = new Mock<IWorkflowActivator>();
var workflowActivatorMock = new Mock<IWorkflowFactory>();
var workflowRegistryMock = new Mock<IWorkflowRegistry>();
var workflowInstanceManager = new WorkflowInstanceManager(_session);
var workflowExpressionEvaluatorMock = new Mock<IExpressionEvaluator>();
@ -43,8 +43,8 @@ namespace Elsa.Core.UnitTests
var serviceProvider = new ServiceCollection().BuildServiceProvider();
workflowActivatorMock
.Setup(x => x.ActivateAsync(It.IsAny<Workflow>(), It.IsAny<string?>(), It.IsAny<CancellationToken>()))
.ReturnsAsync((Workflow workflow, string? correlationId, CancellationToken cancellationToken) => new WorkflowInstance());
.Setup(x => x.InstantiateAsync(It.IsAny<WorkflowBlueprint>(), It.IsAny<string?>(), It.IsAny<CancellationToken>()))
.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;
}