Incremental work on jobs as activities

This commit is contained in:
Sipke Schoorstra 2022-08-10 16:02:53 +02:00
parent 3174399517
commit dd217c7e8b
37 changed files with 292 additions and 103 deletions

View file

@ -1,14 +0,0 @@
using Elsa.Jobs.Models;
using Elsa.Jobs.Services;
namespace Elsa.WorkflowServer.Web;
public class IndexBlockchainJob : IJob
{
public async ValueTask ExecuteAsync(JobExecutionContext context)
{
Console.WriteLine("Indexing blockchain...");
await Task.Delay(5000);
Console.WriteLine("Finished indexing blockchain.");
}
}

View file

@ -0,0 +1,20 @@
using Elsa.Activities.Jobs.Features;
using Elsa.Jobs.Abstractions;
using Elsa.Jobs.Models;
using Elsa.Jobs.Services;
namespace Elsa.WorkflowServer.Web.Jobs;
/// <summary>
/// Jobs can be scheduled manually using <see cref="IJobQueue"/>,
/// but when enabling the <see cref="JobsFeature"/>, these jobs become available as activities too.
/// </summary>
public class IndexBlockchainJob : Job
{
protected override async ValueTask ExecuteAsync(JobExecutionContext context)
{
Console.WriteLine("Indexing blockchain...");
await Task.Delay(5000);
Console.WriteLine("Finished indexing blockchain.");
}
}

View file

@ -37,6 +37,7 @@ using Elsa.Workflows.Persistence.Extensions;
using Elsa.Workflows.Runtime.Extensions;
using Elsa.WorkflowServer.Web;
using Elsa.WorkflowServer.Web.Implementations;
using Elsa.WorkflowServer.Web.Jobs;
using FastEndpoints;
using FastEndpoints.Security;

View file

@ -6,6 +6,8 @@ namespace Elsa.Jobs.Abstractions;
public abstract class Job : IJob
{
public string Id { get; set; } = default!;
ValueTask IJob.ExecuteAsync(JobExecutionContext context) => ExecuteAsync(context);
protected virtual ValueTask ExecuteAsync(JobExecutionContext context)

View file

@ -11,4 +11,8 @@
<PackageReference Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="6.0.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\modules\Elsa.Mediator\Elsa.Mediator.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,11 @@
using System;
namespace Elsa.Jobs.Services;
public class JobFactory : IJobFactory
{
public IJob Create(Type jobType)
{
throw new NotImplementedException();
}
}

View file

@ -2,16 +2,20 @@ using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Jobs.Models;
using Elsa.Jobs.Notifications;
using Elsa.Jobs.Services;
using Elsa.Mediator.Services;
namespace Elsa.Jobs.Implementations;
public class JobRunner : IJobRunner
{
private readonly IEventPublisher _eventPublisher;
private readonly IServiceProvider _serviceProvider;
public JobRunner(IServiceProvider serviceProvider)
public JobRunner(IEventPublisher eventPublisher, IServiceProvider serviceProvider)
{
_eventPublisher = eventPublisher;
_serviceProvider = serviceProvider;
}
@ -19,5 +23,6 @@ public class JobRunner : IJobRunner
{
var context = new JobExecutionContext(_serviceProvider, cancellationToken);
await job.ExecuteAsync(context);
await _eventPublisher.PublishAsync(new JobExecuted(job), cancellationToken);
}
}

View file

@ -0,0 +1,6 @@
using Elsa.Jobs.Services;
using Elsa.Mediator.Services;
namespace Elsa.Jobs.Notifications;
public record JobExecuted(IJob Job) : INotification;

View file

@ -5,5 +5,6 @@ namespace Elsa.Jobs.Services;
public interface IJob
{
string Id { get; set; }
ValueTask ExecuteAsync(JobExecutionContext context);
}

View file

@ -0,0 +1,11 @@
using System;
namespace Elsa.Jobs.Services;
/// <summary>
/// Instantiates new jobs of a given type.
/// </summary>
public interface IJobFactory
{
IJob Create(Type jobType);
}

View file

@ -8,5 +8,5 @@ namespace Elsa.Jobs.Services;
/// </summary>
public interface IJobQueue
{
Task SubmitJobAsync(IJob job, string? queueName = default, CancellationToken cancellationToken = default);
Task<string> SubmitJobAsync(IJob job, string? queueName = default, CancellationToken cancellationToken = default);
}

View file

@ -38,6 +38,7 @@ export class Canvas {
@Method()
public async addActivity(args: AddActivityArgs): Promise<Activity> {
console.debug(`Activity type: ${args.descriptor.activityType}`)
return await this.root.addActivity(args);
}

View file

@ -17,9 +17,7 @@ interface ActivityCategoryModel {
})
export class ToolboxActivities {
@Prop() graph: Graph;
//@State() activityCategoryModels: Array<ActivityCategoryModel> = [];
private dnd: Addon.Dnd;
//private renderedActivities: Map<string, string>;
@State() private expandedCategories: Array<string> = [];
@Watch('graph')
@ -35,10 +33,6 @@ export class ToolboxActivities {
});
}
componentWillLoad() {
//this.handleActivityDescriptorsChanged(descriptorsStore.activityDescriptors);
}
private static onActivityStartDrag(e: DragEvent, activityDescriptor: ActivityDescriptor) {
const json = JSON.stringify(activityDescriptor);
e.dataTransfer.setData('activity-descriptor', json);
@ -55,37 +49,6 @@ export class ToolboxActivities {
this.expandedCategories = [...expandedCategories, category];
}
// handleActivityDescriptorsChanged(value: Array<ActivityDescriptor>) {
// const browsableDescriptors = value.filter(x => x.isBrowsable);
// const categorizedActivitiesLookup = groupBy(browsableDescriptors, x => x.category);
// const categories = Object.keys(categorizedActivitiesLookup);
// const renderedActivities: Map<string, string> = new Map<string, string>();
//
// // Group activities by category
// this.activityCategoryModels = categories.map(x => {
// const model: ActivityCategoryModel = {
// category: x,
// expanded: false,
// activities: categorizedActivitiesLookup[x]
// };
//
// return model;
// });
//
// // Render activities.
// const activityDriverRegistry = Container.get(ActivityDriverRegistry);
//
// for (const activityDescriptor of browsableDescriptors) {
// const activityType = activityDescriptor.activityType;
// const driver = activityDriverRegistry.createDriver(activityType);
// const html = driver.display({displayType: 'picker', activityDescriptor: activityDescriptor});
//
// renderedActivities.set(activityType, html);
// }
//
// this.renderedActivities = renderedActivities;
// }
buildModel = (): any => {
const browsableDescriptors = descriptorsStore.activityDescriptors.filter(x => x.isBrowsable);
const categorizedActivitiesLookup = groupBy(browsableDescriptors, x => x.category);
@ -122,8 +85,6 @@ export class ToolboxActivities {
render() {
// const categoryModels = this.activityCategoryModels;
// const renderedActivities = this.renderedActivities;
const model = this.buildModel();
const categoryModels = model.categories;
const renderedActivities = model.activities;

View file

@ -1,4 +1,6 @@
using Elsa.Workflows.Core.Models;
using Elsa.Activities.Jobs.Models;
using Elsa.Jobs.Services;
using Elsa.Workflows.Core.Models;
namespace Elsa.Activities.Jobs.Activities;
@ -15,6 +17,20 @@ public class JobActivity : ActivityBase
{
JobType = jobType;
}
public Type JobType { get; set; } = default!;
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
var jobQueue = context.GetRequiredService<IJobQueue>();
var job = (IJob)Activator.CreateInstance(JobType)!;
var jobId = await jobQueue.SubmitJobAsync(job, cancellationToken: context.CancellationToken);
var bookmarkPayload = new EnqueuedJobPayload(jobId);
context.CreateBookmark(bookmarkPayload, Resume);
}
private async ValueTask Resume(ActivityExecutionContext context)
{
await CompleteAsync(context);
}
}

View file

@ -0,0 +1,33 @@
namespace Elsa.Activities.Jobs.Attributes;
[AttributeUsage(AttributeTargets.Class)]
public class JobAttribute : Attribute
{
public JobAttribute(string @namespace, string? category, string? description = default)
{
Namespace = @namespace;
Description = description;
Category = category;
}
public JobAttribute(string @namespace, string? description = default)
{
Namespace = @namespace;
Description = description;
Category = @namespace;
}
public JobAttribute(string @namespace, string? typeName, string? description = default, string? category = default)
{
Namespace = @namespace;
TypeName = typeName;
Description = description;
Category = category;
}
public string? Namespace { get; set;}
public string? TypeName { get; set;}
public string? Description { get; set;}
public string? DisplayName { get; set; }
public string? Category { get; set;}
}

View file

@ -10,6 +10,7 @@
<ProjectReference Include="..\..\common\Elsa.Jobs.Abstractions\Elsa.Jobs.Abstractions.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Core\Elsa.Workflows.Core.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Management\Elsa.Workflows.Management.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>
<ItemGroup>

View file

@ -0,0 +1,26 @@
using Elsa.Activities.Jobs.Helpers;
using Elsa.Activities.Jobs.Models;
using Elsa.Jobs.Notifications;
using Elsa.Mediator.Services;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Activities.Jobs.Handlers;
public class JobExecutedHandler : INotificationHandler<JobExecuted>
{
private readonly IWorkflowService _workflowService;
public JobExecutedHandler(IWorkflowService workflowService)
{
_workflowService = workflowService;
}
public async Task HandleAsync(JobExecuted notification, CancellationToken cancellationToken)
{
var payload = new EnqueuedJobPayload(notification.Job.Id);
var jobType = notification.Job.GetType();
var jobTypeName = JobTypeNameHelper.GenerateTypeName(jobType);
var bookmarkName = jobTypeName;
await _workflowService.DispatchStimulusAsync(bookmarkName, payload, cancellationToken: cancellationToken);
}
}

View file

@ -0,0 +1,41 @@
using System.Reflection;
using Elsa.Activities.Jobs.Attributes;
using Elsa.Jobs.Services;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Services;
namespace Elsa.Activities.Jobs.Helpers;
public static class JobTypeNameHelper
{
public static string? GenerateNamespace(Type jobType)
{
var activityAttr = jobType.GetCustomAttribute<JobAttribute>();
return activityAttr?.Namespace ?? jobType.Namespace;
}
public static string GenerateTypeName(Type type, string? ns)
{
var activityAttr = type.GetCustomAttribute<JobAttribute>();
var typeName = activityAttr?.TypeName ?? type.Name;
return ns != null ? $"{ns}.{typeName}" : typeName;
}
public static string GenerateTypeName<T>() where T : IJob => GenerateTypeName(typeof(T));
public static string GenerateTypeName(Type type)
{
var ns = GenerateNamespace(type);
return GenerateTypeName(type, ns);
}
public static string? GetCategoryFromNamespace(string? ns)
{
if (string.IsNullOrWhiteSpace(ns))
return null;
var index = ns.LastIndexOf('.');
return index < 0 ? ns : ns[(index + 1)..];
}
}

View file

@ -1,9 +1,15 @@
using System.ComponentModel;
using System.Reflection;
using Elsa.Activities.Jobs.Activities;
using Elsa.Activities.Jobs.Attributes;
using Elsa.Activities.Jobs.Helpers;
using Elsa.Activities.Jobs.Services;
using Elsa.Jobs.Services;
using Elsa.Workflows.Core.Helpers;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Management.Extensions;
using Elsa.Workflows.Management.Services;
using Humanizer;
namespace Elsa.Activities.Jobs.Implementations;
@ -21,31 +27,40 @@ public class JobActivityProvider : IActivityProvider
_jobRegistry = jobRegistry;
}
public async ValueTask<IEnumerable<ActivityDescriptor>> GetDescriptorsAsync(CancellationToken cancellationToken = default)
public ValueTask<IEnumerable<ActivityDescriptor>> GetDescriptorsAsync(CancellationToken cancellationToken = default)
{
var jobTypes = _jobRegistry.List();
var descriptors = CreateDescriptors(jobTypes).ToList();
return descriptors;
return new(descriptors);
}
private IEnumerable<ActivityDescriptor> CreateDescriptors(IEnumerable<Type> jobTypes) => jobTypes.Select(CreateDescriptor);
private ActivityDescriptor CreateDescriptor(Type jobType)
{
var typeName = jobType.Name;
var jobAttr = jobType.GetCustomAttribute<JobAttribute>();
var ns = jobAttr?.Namespace ?? JobTypeNameHelper.GenerateNamespace(jobType);
var typeName = jobAttr?.TypeName ?? jobType.Name;
var fullTypeName = JobTypeNameHelper.GenerateTypeName(jobType);
var displayNameAttr = jobType.GetCustomAttribute<DisplayNameAttribute>();
var displayName = displayNameAttr?.DisplayName ?? jobAttr?.DisplayName ?? typeName.Humanize(LetterCasing.Title);
var categoryAttr = jobType.GetCustomAttribute<CategoryAttribute>();
var category = categoryAttr?.Category ?? jobAttr?.Category ?? ActivityTypeNameHelper.GetCategoryFromNamespace(ns) ?? "Miscellaneous";
var descriptionAttr = jobType.GetCustomAttribute<DescriptionAttribute>();
var description = descriptionAttr?.Description ?? jobAttr?.Description;
return new()
{
ActivityType = typeName,
DisplayName = jobType.Name,
Description = "",
Category = "Jobs",
ActivityType = fullTypeName,
DisplayName = displayName,
Description = description,
Category = category,
Kind = ActivityKind.Job,
IsBrowsable = true,
Constructor = context =>
{
var activity = _activityFactory.Create<JobActivity>(context);
activity.TypeName = typeName;
activity.TypeName = fullTypeName;
activity.JobType = jobType;
return activity;

View file

@ -0,0 +1,19 @@
using System.Text.Json.Serialization;
namespace Elsa.Activities.Jobs.Models;
public record EnqueuedJobPayload
{
[JsonConstructor]
public EnqueuedJobPayload()
{
}
public EnqueuedJobPayload(string jobId)
{
JobId = jobId;
}
public string JobId { get; init; } = default!;
}

View file

@ -55,7 +55,7 @@ public class MessageReceived : Trigger<object>
/// </summary>
public IFormatter? Formatter { get; set; }
protected override object GetTriggerDatum(TriggerIndexingContext context) => GetBookmarkData(context.ExpressionExecutionContext);
protected override object GetTriggerDatum(TriggerIndexingContext context) => GetBookmarkPayload(context.ExpressionExecutionContext);
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
@ -63,7 +63,7 @@ public class MessageReceived : Trigger<object>
if (!context.TryGetInput<ReceivedServiceBusMessageModel>(InputKey, out var receivedMessage))
{
// Create bookmarks for when we receive the expected HTTP request.
context.CreateBookmark(GetBookmarkData(context.ExpressionExecutionContext), Resume);
context.CreateBookmark(GetBookmarkPayload(context.ExpressionExecutionContext), Resume);
return;
}
@ -87,7 +87,7 @@ public class MessageReceived : Trigger<object>
context.Set(Result, body);
}
private object GetBookmarkData(ExpressionExecutionContext context)
private object GetBookmarkPayload(ExpressionExecutionContext context)
{
var queueOrTopic = context.Get(QueueOrTopic)!;
var subscription = context.Get(Subscription);

View file

@ -14,10 +14,10 @@ namespace Elsa.AzureServiceBus.Handlers;
/// </summary>
public class UpdateWorkers : INotificationHandler<WorkflowTriggersIndexed>, INotificationHandler<WorkflowBookmarksDeleted>, INotificationHandler<WorkflowBookmarksSaved>
{
private readonly IBookmarkDataSerializer _serializer;
private readonly IBookmarkPayloadSerializer _serializer;
private readonly IWorkerManager _workerManager;
public UpdateWorkers(IBookmarkDataSerializer serializer, IWorkerManager workerManager)
public UpdateWorkers(IBookmarkPayloadSerializer serializer, IWorkerManager workerManager)
{
_serializer = serializer;
_workerManager = workerManager;

View file

@ -71,12 +71,10 @@ public class Worker : IAsyncDisposable
private async Task InvokeWorkflowsAsync(ServiceBusReceivedMessage message, CancellationToken cancellationToken)
{
var payload = new MessageReceivedTriggerPayload(QueueOrTopic, Subscription);
var hash = _hasher.Hash(payload);
var correlationId = message.CorrelationId;
var messageModel = CreateMessageModel(message);
var input = new Dictionary<string, object>() { [MessageReceived.InputKey] = messageModel };
var stimulus = Stimulus.Standard(BookmarkName, hash, input, correlationId);
var executionResults = (await _workflowService.DispatchStimulusAsync(stimulus, cancellationToken)).ToList();
var input = new Dictionary<string, object> { [MessageReceived.InputKey] = messageModel };
var executionResults = (await _workflowService.DispatchStimulusAsync(BookmarkName, payload, input, correlationId, cancellationToken)).ToList();
_logger.LogInformation("Triggered {WorkflowCount} workflows", executionResults.Count);
}

View file

@ -1,6 +1,7 @@
using Elsa.Hangfire.Jobs;
using Elsa.Jobs.Services;
using Hangfire;
using Hangfire.Server;
using Hangfire.States;
using HangfireJob = Hangfire.Common.Job;
@ -10,16 +11,16 @@ public class HangfireJobQueue : IJobQueue
{
private readonly IBackgroundJobClient _backgroundJobClient;
public HangfireJobQueue(IBackgroundJobClient backgroundJobClient, IJobRunner jobRunner)
public HangfireJobQueue(IBackgroundJobClient backgroundJobClient)
{
_backgroundJobClient = backgroundJobClient;
}
public Task SubmitJobAsync(IJob job, string? queueName = default, CancellationToken cancellationToken = default)
public Task<string> SubmitJobAsync(IJob job, string? queueName = default, CancellationToken cancellationToken = default)
{
var hangfireJob = HangfireJob.FromExpression<RunElsaJob>(x => x.RunAsync(job, CancellationToken.None));
_backgroundJobClient.Create(hangfireJob, new EnqueuedState(queueName ?? "default"));
var jobId = _backgroundJobClient.Create(hangfireJob, new EnqueuedState(queueName ?? "default"));
return Task.CompletedTask;
return Task.FromResult(jobId);
}
}

View file

@ -1,4 +1,5 @@
using Elsa.Jobs.Services;
using Hangfire.Server;
namespace Elsa.Hangfire.Jobs;

View file

@ -41,7 +41,7 @@ public class HttpEndpoint : Trigger<HttpRequestModel>
)]
public Input<string?> Policy { get; set; } = new(default(string?));
protected override IEnumerable<object> GetTriggerData(TriggerIndexingContext context) => GetBookmarkData(context.ExpressionExecutionContext);
protected override IEnumerable<object> GetTriggerPayload(TriggerIndexingContext context) => GetBookmarkPayload(context.ExpressionExecutionContext);
protected override void Execute(ActivityExecutionContext context)
{
@ -49,7 +49,7 @@ public class HttpEndpoint : Trigger<HttpRequestModel>
if (!context.TryGetInput<HttpRequestModel>(InputKey, out var request))
{
// Create bookmarks for when we receive the expected HTTP request.
context.CreateBookmarks(GetBookmarkData(context.ExpressionExecutionContext));
context.CreateBookmarks(GetBookmarkPayload(context.ExpressionExecutionContext));
return;
}
@ -57,7 +57,7 @@ public class HttpEndpoint : Trigger<HttpRequestModel>
context.Set(Result, request);
}
private IEnumerable<object> GetBookmarkData(ExpressionExecutionContext context)
private IEnumerable<object> GetBookmarkPayload(ExpressionExecutionContext context)
{
// Generate bookmark data for path and selected methods.
var path = context.Get(Path);

View file

@ -41,7 +41,7 @@ public class StartAt : Trigger
[Input] public Input<DateTimeOffset> DateTime { get; set; } = default!;
protected override object GetTriggerDatum(TriggerIndexingContext context)
protected override object GetTriggerPayload(TriggerIndexingContext context)
{
var executeAt = context.ExpressionExecutionContext.Get(DateTime);
return new StartAtPayload(executeAt);

View file

@ -29,7 +29,7 @@ public class Timer : EventGenerator
[Input] public Input<TimeSpan> Interval { get; set; } = default!;
protected override object GetTriggerDatum(TriggerIndexingContext context)
protected override object GetTriggerPayload(TriggerIndexingContext context)
{
var interval = context.ExpressionExecutionContext.Get(Interval);
var clock = context.ExpressionExecutionContext.GetRequiredService<ISystemClock>();

View file

@ -65,7 +65,7 @@ public class WorkflowsFeature : FeatureBase
.AddSingleton<IHasher, Hasher>()
.AddSingleton<IIdentityGenerator, RandomIdentityGenerator>()
.AddSingleton<ISystemClock, SystemClock>()
.AddSingleton<IBookmarkDataSerializer, BookmarkDataSerializer>()
.AddSingleton<IBookmarkPayloadSerializer, BookmarkPayloadSerializer>()
.AddTransient<WorkflowDefinitionBuilder>()
.AddSingleton(typeof(Func<IWorkflowDefinitionBuilder>), sp => () => sp.GetRequiredService<WorkflowDefinitionBuilder>())
.AddSingleton<IWorkflowDefinitionBuilderFactory, WorkflowDefinitionBuilderFactory>()

View file

@ -3,11 +3,11 @@ using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Implementations;
public class BookmarkDataSerializer : IBookmarkDataSerializer
public class BookmarkPayloadSerializer : IBookmarkPayloadSerializer
{
private readonly JsonSerializerOptions _settings;
public BookmarkDataSerializer()
public BookmarkPayloadSerializer()
{
_settings = new JsonSerializerOptions
{

View file

@ -6,16 +6,16 @@ namespace Elsa.Workflows.Core.Implementations;
public class Hasher : IHasher
{
private readonly IBookmarkDataSerializer _bookmarkDataSerializer;
private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer;
public Hasher(IBookmarkDataSerializer bookmarkDataSerializer)
public Hasher(IBookmarkPayloadSerializer bookmarkPayloadSerializer)
{
_bookmarkDataSerializer = bookmarkDataSerializer;
_bookmarkPayloadSerializer = bookmarkPayloadSerializer;
}
public string Hash(object value)
{
var json = _bookmarkDataSerializer.Serialize(value);
var json = _bookmarkPayloadSerializer.Serialize(value);
return Hash(json);
}

View file

@ -111,19 +111,19 @@ public class ActivityExecutionContext
public Bookmark CreateBookmark(ExecuteActivityDelegate callback) => CreateBookmark(default, callback);
public Bookmark CreateBookmark(object? bookmarkDatum = default, ExecuteActivityDelegate? callback = default)
public Bookmark CreateBookmark(object? payload = default, ExecuteActivityDelegate? callback = default)
{
var hasher = GetRequiredService<IHasher>();
var identityGenerator = GetRequiredService<IIdentityGenerator>();
var bookmarkDataSerializer = GetRequiredService<IBookmarkDataSerializer>();
var bookmarkDatumJson = bookmarkDatum != null ? bookmarkDataSerializer.Serialize(bookmarkDatum) : default;
var hash = bookmarkDatumJson != null ? hasher.Hash(bookmarkDatumJson) : default;
var payloadSerializer = GetRequiredService<IBookmarkPayloadSerializer>();
var payloadJson = payload != null ? payloadSerializer.Serialize(payload) : default;
var hash = payloadJson != null ? hasher.Hash(payloadJson) : default;
var bookmark = new Bookmark(
identityGenerator.GenerateId(),
Activity.TypeName,
hash,
bookmarkDatumJson,
payloadJson,
Activity.Id,
Id,
callback?.Method.Name);

View file

@ -26,12 +26,12 @@ public abstract class Trigger : Activity, ITrigger
/// <summary>
/// Override this method to return trigger data.
/// </summary>
protected virtual IEnumerable<object> GetTriggerData(TriggerIndexingContext context) => new[] { GetTriggerDatum(context) };
protected virtual IEnumerable<object> GetTriggerData(TriggerIndexingContext context) => new[] { GetTriggerPayload(context) };
/// <summary>
/// Override this method to return a trigger datum.
/// </summary>
protected virtual object GetTriggerDatum(TriggerIndexingContext context) => new();
protected virtual object GetTriggerPayload(TriggerIndexingContext context) => new();
}
public abstract class Trigger<T> : Activity<T>, ITrigger
@ -55,14 +55,14 @@ public abstract class Trigger<T> : Activity<T>, ITrigger
/// </summary>
protected virtual ValueTask<IEnumerable<object>> GetTriggerDataAsync(TriggerIndexingContext context)
{
var hashes = GetTriggerData(context);
var hashes = GetTriggerPayload(context);
return ValueTask.FromResult(hashes);
}
/// <summary>
/// Override this method to return trigger data.
/// </summary>
protected virtual IEnumerable<object> GetTriggerData(TriggerIndexingContext context) => new[] { GetTriggerDatum(context) };
protected virtual IEnumerable<object> GetTriggerPayload(TriggerIndexingContext context) => new[] { GetTriggerDatum(context) };
/// <summary>
/// Override this method to return a trigger datum.

View file

@ -1,6 +1,6 @@
namespace Elsa.Workflows.Core.Services;
public interface IBookmarkDataSerializer
public interface IBookmarkPayloadSerializer
{
T Deserialize<T>(string json) where T : notnull;
string Serialize<T>(T payload) where T : notnull;

View file

@ -1,5 +1,6 @@
using Elsa.Persistence.Common.Models;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Persistence.Models;
using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Services;
@ -12,17 +13,20 @@ public class WorkflowService : IWorkflowService
private readonly IWorkflowDispatcher _workflowDispatcher;
private readonly IWorkflowInstructionExecutor _workflowInstructionExecutor;
private readonly IStimulusInterpreter _stimulusInterpreter;
private readonly IHasher _hasher;
public WorkflowService(
IWorkflowInvoker workflowInvoker,
IWorkflowDispatcher workflowDispatcher,
IWorkflowInstructionExecutor workflowInstructionExecutor,
IStimulusInterpreter stimulusInterpreter)
IStimulusInterpreter stimulusInterpreter,
IHasher hasher)
{
_workflowInvoker = workflowInvoker;
_workflowDispatcher = workflowDispatcher;
_workflowInstructionExecutor = workflowInstructionExecutor;
_stimulusInterpreter = stimulusInterpreter;
_hasher = hasher;
}
public async Task<InvokeWorkflowResult> ExecuteWorkflowAsync(string definitionId, VersionOptions versionOptions, IDictionary<string, object>? input = default, string? correlationId = default, CancellationToken cancellationToken = default)
@ -57,6 +61,27 @@ public class WorkflowService : IWorkflowService
// Execute instructions.
return await _workflowInstructionExecutor.ExecuteInstructionsAsync(instructions, cancellationToken);
}
public async Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, object? inputs = default, string? correlationId = default, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(bookmarkPayload);
var stimulus = Stimulus.Standard(bookmarkName, hash, inputs, correlationId);
return await DispatchStimulusAsync(stimulus, cancellationToken);
}
public async Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, IDictionary<string, object> inputs, string? correlationId = default, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(bookmarkPayload);
var stimulus = Stimulus.Standard(bookmarkName, hash, inputs, correlationId);
return await DispatchStimulusAsync(stimulus, cancellationToken);
}
public async Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, string? correlationId = default, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(bookmarkPayload);
var stimulus = Stimulus.Standard(bookmarkName, hash, default, correlationId);
return await DispatchStimulusAsync(stimulus, cancellationToken);
}
public async Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(IStimulus stimulus, CancellationToken cancellationToken)
{

View file

@ -17,6 +17,6 @@ public static class Stimulus
new(ActivityTypeNameHelper.GenerateTypeName<T>(), default, input.ToDictionary(), correlationId);
public static StandardStimulus Standard(string activityTypeName, string? hash = default, IDictionary<string, object>? input = default, string? correlationId = default) => new(activityTypeName, hash, input, correlationId);
public static StandardStimulus Standard(string activityTypeName, string? hash, object input, string? correlationId = default) => new(activityTypeName, hash, input.ToDictionary(), correlationId);
public static StandardStimulus Standard(string activityTypeName, object input, string? correlationId = default) => new(activityTypeName, default, input.ToDictionary(), correlationId);
public static StandardStimulus Standard(string activityTypeName, string? hash, object? input, string? correlationId = default) => new(activityTypeName, hash, input?.ToDictionary(), correlationId);
public static StandardStimulus Standard(string activityTypeName, object? input, string? correlationId = default) => new(activityTypeName, default, input?.ToDictionary(), correlationId);
}

View file

@ -1,6 +1,7 @@
using Elsa.Persistence.Common.Models;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Persistence.Models;
using Elsa.Workflows.Runtime.Attributes;
using Elsa.Workflows.Runtime.Models;
namespace Elsa.Workflows.Runtime.Services;
@ -16,5 +17,8 @@ public interface IWorkflowService
Task<DispatchWorkflowInstanceResponse> DispatchWorkflowAsync(string instanceId, Bookmark bookmark, IDictionary<string, object>? input = default, string? correlationId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<ExecuteWorkflowInstructionResult>> ExecuteStimulusAsync(IStimulus stimulus, CancellationToken cancellationToken = default);
Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(IStimulus stimulus, CancellationToken cancellationToken = default);
Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, object inputs, string? correlationId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, IDictionary<string, object> inputs, string? correlationId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<DispatchWorkflowInstructionResult>> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, string? correlationId = default, CancellationToken cancellationToken = default);
}