Complete workflow context sample and fix NodaTime serialization

This commit is contained in:
Sipke Schoorstra 2020-10-31 22:19:58 +01:00
parent ef660ab541
commit ce6509baff
29 changed files with 197 additions and 83 deletions

View file

@ -15,5 +15,8 @@ namespace Elsa.Activities.Http
public static IActivityBuilder ReceiveHttpRequest(this IBuilder builder, Func<ValueTask<PathString>> path) => builder.ReceiveHttpRequest(setup => setup.Set(x => x.Path, path));
public static IActivityBuilder ReceiveHttpRequest(this IBuilder builder, Func<PathString> path) => builder.ReceiveHttpRequest(setup => setup.Set(x => x.Path, path));
public static IActivityBuilder ReceiveHttpRequest(this IBuilder builder, PathString path) => builder.ReceiveHttpRequest(setup => setup.Set(x => x.Path, path));
public static IActivityBuilder ReceiveHttpPostRequest<T>(this IBuilder builder, Func<ActivityExecutionContext, PathString> path) => builder.ReceiveHttpRequest(activity => activity.WithPath(path).WithMethod(HttpMethods.Post).WithTargetType<T>());
public static IActivityBuilder ReceiveHttpPostRequest<T>(this IBuilder builder, Func<PathString> path) => builder.ReceiveHttpRequest(activity => activity.WithPath(path).WithMethod(HttpMethods.Post).WithTargetType<T>());
public static IActivityBuilder ReceiveHttpPostRequest<T>(this IBuilder builder, PathString path) => builder.ReceiveHttpRequest(activity => activity.WithPath(path).WithMethod(HttpMethods.Post).WithTargetType<T>());
}
}

View file

@ -1,6 +1,5 @@
using System;
using System.IO;
using System.Linq;
using System.Text;
using System.Threading;
using System.Threading.Tasks;

View file

@ -1,4 +1,3 @@
using System;
using System.Linq;
using System.Net;
using System.Threading.Tasks;

View file

@ -4,13 +4,20 @@ using System.Threading;
using System.Threading.Tasks;
using Elsa.Activities.Http.Extensions;
using Elsa.Activities.Http.Services;
using Elsa.Serialization;
using Microsoft.AspNetCore.Http;
using Newtonsoft.Json;
namespace Elsa.Activities.Http.Parsers
{
public class JsonHttpRequestBodyParser : IHttpRequestBodyParser
{
private readonly IContentSerializer _serializer;
public JsonHttpRequestBodyParser(IContentSerializer serializer)
{
_serializer = serializer;
}
public int Priority => 0;
public string?[] SupportedContentTypes => new[] { "application/json", "text/json" };
@ -18,7 +25,7 @@ namespace Elsa.Activities.Http.Parsers
{
var json = await request.ReadContentAsStringAsync(cancellationToken);
targetType ??= typeof(ExpandoObject);
return JsonConvert.DeserializeObject(json, targetType)!;
return _serializer.Deserialize(json, targetType)!;
}
}
}

View file

@ -2,8 +2,6 @@ using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
using NCrontab;
using NodaTime;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.Timers

View file

@ -3,7 +3,6 @@ using Elsa.Activities.Timers;
using Elsa.Activities.Timers.HostedServices;
using Elsa.Activities.Timers.Options;
using Elsa.Activities.Timers.Triggers;
using Elsa.Triggers;
// ReSharper disable once CheckNamespace
namespace Microsoft.Extensions.DependencyInjection

View file

@ -1,12 +1,9 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Extensions;
using Elsa.Indexes;
using Elsa.Models;
using Elsa.Services;
using Elsa.Triggers;
using NodaTime;
using Open.Linq.AsyncExtensions;
namespace Elsa.Activities.Timers.Triggers
{

View file

@ -19,6 +19,7 @@ namespace Elsa.Builders
IActivityBuilder Add<T>(Action<ISetupActivity<T>>? setup = default) where T : class, IActivity;
IOutcomeBuilder When(string outcome);
IActivityBuilder Then(IActivityBuilder targetActivity);
IConnectionBuilder Then(string activityName);
IActivityBuilder WithId(string? id);
IActivityBuilder WithName(string? name);
Func<ActivityExecutionContext, CancellationToken, ValueTask<IActivity>> BuildActivityAsync();

View file

@ -13,5 +13,6 @@ namespace Microsoft.Extensions.DependencyInjection
public static IServiceCollection AddDataMigration<T>(this IServiceCollection services) where T : class, IDataMigration => services.AddScoped<IDataMigration, T>();
public static IServiceCollection AddWorkflowProvider<T>(this IServiceCollection services) where T : class, IWorkflowProvider => services.AddTransient<IWorkflowProvider, T>();
public static IServiceCollection AddTriggerProvider<T>(this IServiceCollection services) where T : class, ITriggerProvider => services.AddTransient<ITriggerProvider, T>();
public static IServiceCollection AddWorkflowContextProvider<T>(this IServiceCollection services) where T : class, IWorkflowContextProvider => services.AddTransient<IWorkflowContextProvider, T>();
}
}

View file

@ -1,4 +1,5 @@
using Newtonsoft.Json.Linq;
using System;
using Newtonsoft.Json.Linq;
namespace Elsa.Serialization
{
@ -6,7 +7,9 @@ namespace Elsa.Serialization
{
string Serialize<T>(T value);
T Deserialize<T>(JToken token);
object? Deserialize(JToken token, Type targetType);
T Deserialize<T>(string json);
object? Deserialize(string json, Type targetType);
object GetSettings();
}
}

View file

@ -2,7 +2,6 @@ using System;
using System.Collections.Generic;
using System.Linq;
using Elsa.Models;
using Elsa.Serialization;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Localization;
using Newtonsoft.Json;

View file

@ -0,0 +1,15 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Services.Models;
namespace Elsa.Services
{
public abstract class WorkflowContextProvider<T> : IWorkflowContextProvider
{
public IEnumerable<Type> SupportedTypes => new[] { typeof(T) };
public virtual ValueTask<object?> LoadContextAsync(LoadWorkflowContext context, CancellationToken cancellationToken = default) => new ValueTask<object?>();
public virtual ValueTask<string?> SaveContextAsync(SaveWorkflowContext context, CancellationToken cancellationToken = default) => new ValueTask<string?>();
}
}

View file

@ -1,5 +1,4 @@
using Elsa.Models;
using Elsa.Services.Models;
using Elsa.Services.Models;
namespace Elsa.Triggers
{

View file

@ -1,10 +0,0 @@
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
public static class JoinBuilderExtensions
{
}
}

View file

@ -0,0 +1,14 @@
using System;
using Elsa.Builders;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
public static class JoinExtensions
{
public static ISetupActivity<Join> WithMode(this ISetupActivity<Join> activity, Func<ActivityExecutionContext, Join.JoinMode> value) => activity.Set(x => x.Mode, value);
public static ISetupActivity<Join> WithMode(this ISetupActivity<Join> activity, Func<Join.JoinMode> value) => activity.Set(x => x.Mode, value);
public static ISetupActivity<Join> WithMode(this ISetupActivity<Join> activity, Join.JoinMode value) => activity.Set(x => x.Mode, value);
}
}

View file

@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Services;
@ -29,17 +30,20 @@ namespace Elsa.Builders
public IActivityBuilder Add<T>(
Action<ISetupActivity<T>>? setup = default)
where T : class, IActivity => WorkflowBuilder.Add(setup);
where T : class, IActivity =>
WorkflowBuilder.Add(setup);
public IOutcomeBuilder When(string outcome) => new OutcomeBuilder(WorkflowBuilder, this, outcome);
public IActivityBuilder Then<T>(
Action<ISetupActivity<T>>? setup = null,
Action<IActivityBuilder>? branch = null)
where T : class, IActivity => When(OutcomeNames.Done).Then(setup, branch);
where T : class, IActivity =>
When(OutcomeNames.Done).Then(setup, branch);
public IActivityBuilder Then<T>(Action<IActivityBuilder>? branch = null)
where T : class, IActivity => When(OutcomeNames.Done).Then<T>(branch);
where T : class, IActivity =>
When(OutcomeNames.Done).Then<T>(branch);
public IActivityBuilder Then(IActivityBuilder targetActivity)
{
@ -47,6 +51,11 @@ namespace Elsa.Builders
return this;
}
public IConnectionBuilder Then(string activityName) =>
WorkflowBuilder.Connect(
() => this,
() => WorkflowBuilder.Activities.First(x => x.Name == activityName));
public IActivityBuilder WithId(string? id)
{
ActivityId = id!;

View file

@ -1,6 +1,5 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Data;
using Elsa.Services;
// ReSharper disable once CheckNamespace

View file

@ -1,3 +1,4 @@
using System;
using Elsa.Converters;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
@ -12,19 +13,17 @@ namespace Elsa.Serialization
public DefaultContentSerializer(JsonSerializer serializer)
{
Serializer = serializer;
SerializerSettings = CreateDefaultJsonSerializationSettings();
}
private JsonSerializerSettings SerializerSettings { get; }
private JsonSerializer Serializer { get; }
public string Serialize<T>(T value) => JObject.FromObject(value!, Serializer).ToString();
public T Deserialize<T>(JToken token) => token.ToObject<T>(Serializer)!;
public T Deserialize<T>(string json)
{
var token = JObject.Parse(json);
return Deserialize<T>(token);
}
public object GetSettings() => CreateDefaultJsonSerializationSettings();
public object? Deserialize(JToken token, Type targetType) => token.ToObject(targetType, Serializer);
public T Deserialize<T>(string json) => JsonConvert.DeserializeObject<T>(json, SerializerSettings)!;
public object? Deserialize(string json, Type targetType) => JsonConvert.DeserializeObject(json, targetType, SerializerSettings);
public object GetSettings() => SerializerSettings;
public static void ConfigureDefaultJsonSerializationSettings(JsonSerializerSettings settings)
{

View file

@ -18,6 +18,7 @@ namespace Elsa.Triggers
private readonly IWorkflowRegistry _workflowRegistry;
private readonly IWorkflowFactory _workflowFactory;
private readonly IWorkflowInstanceManager _workflowInstanceManager;
private readonly IWorkflowContextManager _workflowContextManager;
private readonly IEnumerable<ITriggerProvider> _triggerProviders;
private readonly IMemoryCache _memoryCache;
private readonly IServiceProvider _serviceProvider;
@ -27,6 +28,7 @@ namespace Elsa.Triggers
IWorkflowRegistry workflowRegistry,
IWorkflowFactory workflowFactory,
IWorkflowInstanceManager workflowInstanceManager,
IWorkflowContextManager workflowContextManager,
IEnumerable<ITriggerProvider> triggerProviders,
IMemoryCache memoryCache,
IServiceProvider serviceProvider)
@ -34,6 +36,7 @@ namespace Elsa.Triggers
_workflowRegistry = workflowRegistry;
_workflowFactory = workflowFactory;
_workflowInstanceManager = workflowInstanceManager;
_workflowContextManager = workflowContextManager;
_triggerProviders = triggerProviders;
_memoryCache = memoryCache;
_serviceProvider = serviceProvider;
@ -152,7 +155,9 @@ namespace Elsa.Triggers
{
var providers = _triggerProviders.ToList();
var descriptors = new List<TriggerDescriptor>();
var workflowExecutionContext = new WorkflowExecutionContext(_serviceProvider, workflowBlueprint, workflowInstance, default, default);
var loadWorkflowContext = new LoadWorkflowContext(workflowBlueprint, workflowInstance);
var workflowContext = await _workflowContextManager.LoadContext(loadWorkflowContext, cancellationToken);
var workflowExecutionContext = new WorkflowExecutionContext(_serviceProvider, workflowBlueprint, workflowInstance, default, workflowContext);
foreach (var blockingActivity in blockingActivities)
{

View file

@ -2,6 +2,7 @@
<PropertyGroup>
<TargetFramework>netcoreapp3.1</TargetFramework>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>

View file

@ -0,0 +1,22 @@
using Elsa.Samples.ContextualWorkflowHttp.Models;
using YesSql.Indexes;
namespace Elsa.Samples.ContextualWorkflowHttp.Indexes
{
public class DocumentIndex : MapIndex
{
public string DocumentUid { get; set; } = default!; // DocumentId is a reserved column name by YesSql, so taking DocumentUid instead.
}
public class DocumentIndexProvider : IndexProvider<Document>
{
public override void Describe(DescribeContext<Document> context)
{
context.For<DocumentIndex>().Map(
x => new DocumentIndex
{
DocumentUid = x.DocumentId
});
}
}
}

View file

@ -0,0 +1,15 @@
using Elsa.Data;
using Elsa.Samples.ContextualWorkflowHttp.Indexes;
using YesSql.Sql;
namespace Elsa.Samples.ContextualWorkflowHttp
{
public class Migrations : DataMigration
{
public int Create()
{
SchemaBuilder.CreateMapIndexTable<DocumentIndex>(table => table.Column<string>(nameof(DocumentIndex.DocumentUid)));
return 1;
}
}
}

View file

@ -4,9 +4,10 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Models
{
public class Document
{
public string Id { get; set; }
public string Title { get; set; }
public string Body { get; set; }
public int Id { get; set; }
public string DocumentId { get; set; } = default!;
public string Title { get; set; } = default!;
public string Body { get; set; } = default!;
public ICollection<Comment> Comments { get; set; } = new List<Comment>();
}
}

View file

@ -6,7 +6,7 @@
"environmentVariables": {
"ASPNETCORE_ENVIRONMENT": "Development"
},
"applicationUrl": "http://localhost:8201"
"applicationUrl": "http://localhost:7301"
}
}
}

View file

@ -1,4 +1,6 @@
using System.Data;
using Elsa.Samples.ContextualWorkflowHttp.Indexes;
using Elsa.Samples.ContextualWorkflowHttp.WorkflowContextProviders;
using Elsa.Samples.ContextualWorkflowHttp.Workflows;
using Microsoft.AspNetCore.Builder;
using Microsoft.Extensions.DependencyInjection;
@ -14,6 +16,9 @@ namespace Elsa.Samples.ContextualWorkflowHttp
.AddElsa(option => option.UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared", IsolationLevel.ReadUncommitted)))
.AddHttpActivities()
.AddConsoleActivities()
.AddDataMigration<Migrations>()
.AddIndexProvider<DocumentIndexProvider>()
.AddWorkflowContextProvider<DocumentWorkflowContextProvider>()
.AddWorkflow<DocumentApprovalWorkflow>();
}

View file

@ -0,0 +1,37 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Samples.ContextualWorkflowHttp.Indexes;
using Elsa.Services;
using Elsa.Services.Models;
using YesSql;
using Document = Elsa.Samples.ContextualWorkflowHttp.Models.Document;
using IIdGenerator = Elsa.Services.IIdGenerator;
namespace Elsa.Samples.ContextualWorkflowHttp.WorkflowContextProviders
{
public class DocumentWorkflowContextProvider : WorkflowContextProvider<Document>
{
private readonly ISession _session;
private readonly IIdGenerator _idGenerator;
public DocumentWorkflowContextProvider(ISession session, IIdGenerator idGenerator)
{
_session = session;
_idGenerator = idGenerator;
}
public override async ValueTask<object?> LoadContextAsync(LoadWorkflowContext context, CancellationToken cancellationToken = default) =>
await _session.Query<Document, DocumentIndex>(x => x.DocumentUid == context.ContextId).FirstOrDefaultAsync();
public override ValueTask<string?> SaveContextAsync(SaveWorkflowContext context, CancellationToken cancellationToken = default)
{
var document = (Document)context.Context;
if (string.IsNullOrWhiteSpace(document.DocumentId))
document.DocumentId = _idGenerator.Generate();
_session.Save(document);
return new ValueTask<string?>(document.DocumentId);
}
}
}

View file

@ -1,49 +1,41 @@
using System;
using System.Net;
using Elsa.Activities.Console;
using Elsa.Activities.ControlFlow;
using Elsa.Activities.Http;
using Elsa.Activities.Http.Models;
using Elsa.Builders;
using Elsa.Samples.ContextualWorkflowHttp.Models;
using Elsa.Services.Models;
using Microsoft.AspNetCore.Http;
namespace Elsa.Samples.ContextualWorkflowHttp.Workflows
{
/// <summary>
/// Demonstrates loading & saving of the document-to-approve, which is the context (or subject) of the workflow.
/// Demonstrates saving & loading of the document-to-approve, which is the context (or subject) of the workflow.
/// </summary>
public class DocumentApprovalWorkflow : IWorkflow
{
public void Build(IWorkflowBuilder workflow)
{
// Demonstrating that we can create activities and connect to them later on by using the activity builder reference.
var join = workflow.Add<Join>(x => x.WithMode(Join.JoinMode.WaitAny));
join.Finish();
workflow
.StartWith<WriteLine>()
// The subject type of this workflow.
// The workflow context type of this workflow.
.WithContextType<Document>()
// Accept HTTP requests to submit new documents.
.ReceiveHttpRequest(activity => activity.WithPath("/documents").WithMethod(HttpMethods.Post).WithTargetType<Document>())
.ReceiveHttpPostRequest<Document>("/documents")
// Store the document as the workflow subject. It will be saved automatically when the workflow gets suspended.
.Then(context => context.WorkflowExecutionContext.WorkflowContext = context.Input)
// Correlate the workflow by document ID.
.Correlate(context => ((Document)context.WorkflowExecutionContext.WorkflowContext)!.Id)
// Store the document as the workflow context. It will be saved automatically when the workflow gets suspended.
.Then(context => context.WorkflowExecutionContext.WorkflowContext = (Document)((HttpRequestModel)context.Input!).Body)
// Write an HTTP response.
.WriteHttpResponse(
activity => activity
.WithStatusCode(HttpStatusCode.OK)
.WithContentType("text/html")
.WithContent(
context =>
{
var document = (Document)context.WorkflowExecutionContext.WorkflowContext;
return $"Document received with ID {document!.Id}! Awaiting Approve or Reject response.";
}))
.WithContent(context => $"Document received with ID {GetDocumentId(context)}! Awaiting Approve or Reject response."))
// Fork execution into two branches: an Approve branch and a Reject branch.
.Then<Fork>(
@ -52,27 +44,30 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows
{
var approveBranch = fork
.When("Approve")
.ReceiveHttpRequest(activity => activity.WithPath("/approve").WithMethod(HttpMethods.Post).WithTargetType<Comment>())
.ReceiveHttpPostRequest<Comment>(context => $"/documents/{GetDocumentId(context)}/approve")
.Then(StoreComment);
var rejectBranch = fork
.When("Reject")
.ReceiveHttpRequest(activity => activity.WithPath("/reject").WithMethod(HttpMethods.Post).WithTargetType<Comment>())
.ReceiveHttpPostRequest<Comment>(context => $"/documents/{GetDocumentId(context)}/reject")
.Then(StoreComment);
WriteResponse(approveBranch, document => $"Thanks for approving document {document!.Id}!").Then("Join");
WriteResponse(rejectBranch, document => $"Thanks for rejecting document {document!.Id}!");
WriteResponse(approveBranch, document => $"Thanks for approving document {document!.DocumentId}!").Then(join);
WriteResponse(rejectBranch, document => $"Thanks for rejecting document {document!.DocumentId}!").Then(join);
});
}
private void StoreComment(ActivityExecutionContext context)
private static Document GetDocument(ActivityExecutionContext context) => (Document)context.WorkflowExecutionContext.WorkflowContext!;
private static string GetDocumentId(ActivityExecutionContext context) => GetDocument(context).DocumentId;
private static void StoreComment(ActivityExecutionContext context)
{
var document = (Document)context.WorkflowExecutionContext.WorkflowContext;
var document = (Document)context.WorkflowExecutionContext.WorkflowContext!;
var comment = (Comment)((HttpRequestModel)context.Input)!.Body;
document!.Comments.Add(comment);
document.Comments.Add(comment);
}
private IActivityBuilder WriteResponse(IActivityBuilder builder, Func<Document, string> html) =>
private static IActivityBuilder WriteResponse(IBuilder builder, Func<Document, string> html) =>
builder.WriteHttpResponse(
activity => activity
.WithStatusCode(HttpStatusCode.OK)
@ -80,7 +75,7 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows
.WithContent(
context =>
{
var document = (Document)context.WorkflowExecutionContext.WorkflowContext;
var document = GetDocument(context);
return html(document);
}));
}

View file

@ -1,28 +1,31 @@
# Register John.
POST http://localhost:8201/register
# Create document.
POST http://localhost:7301/documents
Content-Type: application/json
{
"name": "John von Neumann",
"email": "john@gmail.com"
"documentId": "document-1",
"title": "Document 1",
"body": "john@gmail.com"
}
###
# Register Julia.
POST http://localhost:8201/register
# Approve document.
POST http://localhost:7301/documents/document-1/approve
Content-Type: application/json
{
"name": "Julia Berger",
"email": "julia@gmail.com"
"author": "Jason",
"text": "Great job!",
"timeStamp": "2020-10-31T20:56:49Z"
}
###
# Confirm John's registration
GET http://localhost:8201/confirm?correlation=john@gmail.com
# Reject document.
POST http://localhost:7301/documents/document-1/reject
Content-Type: application/json
###
# Confirm Julia's registration
GET http://localhost:8201/confirm?correlation=julia@gmail.com
Content-Type: application/json
{
"author": "Laura",
"text": "Nice try.",
"timeStamp": "2020-11-01T20:20:20Z"
}

View file

@ -1,4 +1,3 @@
using System;
using Elsa.Activities.Console;
using Elsa.Activities.Timers;
using Elsa.Builders;