diff --git a/src/bundles/Elsa.WorkflowServer.Web/IndexBlockchainJob.cs b/src/bundles/Elsa.WorkflowServer.Web/IndexBlockchainJob.cs deleted file mode 100644 index 4353c20a7..000000000 --- a/src/bundles/Elsa.WorkflowServer.Web/IndexBlockchainJob.cs +++ /dev/null @@ -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."); - } -} \ No newline at end of file diff --git a/src/bundles/Elsa.WorkflowServer.Web/Jobs/IndexBlockchainJob.cs b/src/bundles/Elsa.WorkflowServer.Web/Jobs/IndexBlockchainJob.cs new file mode 100644 index 000000000..d5a015796 --- /dev/null +++ b/src/bundles/Elsa.WorkflowServer.Web/Jobs/IndexBlockchainJob.cs @@ -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; + +/// +/// Jobs can be scheduled manually using , +/// but when enabling the , these jobs become available as activities too. +/// +public class IndexBlockchainJob : Job +{ + protected override async ValueTask ExecuteAsync(JobExecutionContext context) + { + Console.WriteLine("Indexing blockchain..."); + await Task.Delay(5000); + Console.WriteLine("Finished indexing blockchain."); + } +} \ No newline at end of file diff --git a/src/bundles/Elsa.WorkflowServer.Web/Program.cs b/src/bundles/Elsa.WorkflowServer.Web/Program.cs index dc2254ce0..919236969 100644 --- a/src/bundles/Elsa.WorkflowServer.Web/Program.cs +++ b/src/bundles/Elsa.WorkflowServer.Web/Program.cs @@ -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; diff --git a/src/common/Elsa.Jobs.Abstractions/Abstractions/Job.cs b/src/common/Elsa.Jobs.Abstractions/Abstractions/Job.cs index 05d313370..df6230968 100644 --- a/src/common/Elsa.Jobs.Abstractions/Abstractions/Job.cs +++ b/src/common/Elsa.Jobs.Abstractions/Abstractions/Job.cs @@ -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) diff --git a/src/common/Elsa.Jobs.Abstractions/Elsa.Jobs.Abstractions.csproj b/src/common/Elsa.Jobs.Abstractions/Elsa.Jobs.Abstractions.csproj index f1ddfe89e..af9c3eba4 100644 --- a/src/common/Elsa.Jobs.Abstractions/Elsa.Jobs.Abstractions.csproj +++ b/src/common/Elsa.Jobs.Abstractions/Elsa.Jobs.Abstractions.csproj @@ -11,4 +11,8 @@ + + + + diff --git a/src/common/Elsa.Jobs.Abstractions/Implementations/JobFactory.cs b/src/common/Elsa.Jobs.Abstractions/Implementations/JobFactory.cs new file mode 100644 index 000000000..63f99837e --- /dev/null +++ b/src/common/Elsa.Jobs.Abstractions/Implementations/JobFactory.cs @@ -0,0 +1,11 @@ +using System; + +namespace Elsa.Jobs.Services; + +public class JobFactory : IJobFactory +{ + public IJob Create(Type jobType) + { + throw new NotImplementedException(); + } +} \ No newline at end of file diff --git a/src/common/Elsa.Jobs.Abstractions/Implementations/JobRunner.cs b/src/common/Elsa.Jobs.Abstractions/Implementations/JobRunner.cs index ab51ec69d..2dc8c909c 100644 --- a/src/common/Elsa.Jobs.Abstractions/Implementations/JobRunner.cs +++ b/src/common/Elsa.Jobs.Abstractions/Implementations/JobRunner.cs @@ -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); } } \ No newline at end of file diff --git a/src/common/Elsa.Jobs.Abstractions/Notifications/JobExecuted.cs b/src/common/Elsa.Jobs.Abstractions/Notifications/JobExecuted.cs new file mode 100644 index 000000000..db32c9dd8 --- /dev/null +++ b/src/common/Elsa.Jobs.Abstractions/Notifications/JobExecuted.cs @@ -0,0 +1,6 @@ +using Elsa.Jobs.Services; +using Elsa.Mediator.Services; + +namespace Elsa.Jobs.Notifications; + +public record JobExecuted(IJob Job) : INotification; \ No newline at end of file diff --git a/src/common/Elsa.Jobs.Abstractions/Services/IJob.cs b/src/common/Elsa.Jobs.Abstractions/Services/IJob.cs index e7fa5cb16..beb756e5b 100644 --- a/src/common/Elsa.Jobs.Abstractions/Services/IJob.cs +++ b/src/common/Elsa.Jobs.Abstractions/Services/IJob.cs @@ -5,5 +5,6 @@ namespace Elsa.Jobs.Services; public interface IJob { + string Id { get; set; } ValueTask ExecuteAsync(JobExecutionContext context); } \ No newline at end of file diff --git a/src/common/Elsa.Jobs.Abstractions/Services/IJobFactory.cs b/src/common/Elsa.Jobs.Abstractions/Services/IJobFactory.cs new file mode 100644 index 000000000..9d5232ad2 --- /dev/null +++ b/src/common/Elsa.Jobs.Abstractions/Services/IJobFactory.cs @@ -0,0 +1,11 @@ +using System; + +namespace Elsa.Jobs.Services; + +/// +/// Instantiates new jobs of a given type. +/// +public interface IJobFactory +{ + IJob Create(Type jobType); +} \ No newline at end of file diff --git a/src/common/Elsa.Jobs.Abstractions/Services/IJobQueue.cs b/src/common/Elsa.Jobs.Abstractions/Services/IJobQueue.cs index ea5fb5a9e..7fcf0a2b6 100644 --- a/src/common/Elsa.Jobs.Abstractions/Services/IJobQueue.cs +++ b/src/common/Elsa.Jobs.Abstractions/Services/IJobQueue.cs @@ -8,5 +8,5 @@ namespace Elsa.Jobs.Services; /// public interface IJobQueue { - Task SubmitJobAsync(IJob job, string? queueName = default, CancellationToken cancellationToken = default); + Task SubmitJobAsync(IJob job, string? queueName = default, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/designer/elsa-workflows-designer/src/components/designer/canvas/canvas.tsx b/src/designer/elsa-workflows-designer/src/components/designer/canvas/canvas.tsx index 732bbabb2..5b3ce66ee 100644 --- a/src/designer/elsa-workflows-designer/src/components/designer/canvas/canvas.tsx +++ b/src/designer/elsa-workflows-designer/src/components/designer/canvas/canvas.tsx @@ -38,6 +38,7 @@ export class Canvas { @Method() public async addActivity(args: AddActivityArgs): Promise { + console.debug(`Activity type: ${args.descriptor.activityType}`) return await this.root.addActivity(args); } diff --git a/src/designer/elsa-workflows-designer/src/modules/workflow-definitions/components/toolbox-activities.tsx b/src/designer/elsa-workflows-designer/src/modules/workflow-definitions/components/toolbox-activities.tsx index ae893d7c4..84dcf054a 100644 --- a/src/designer/elsa-workflows-designer/src/modules/workflow-definitions/components/toolbox-activities.tsx +++ b/src/designer/elsa-workflows-designer/src/modules/workflow-definitions/components/toolbox-activities.tsx @@ -17,9 +17,7 @@ interface ActivityCategoryModel { }) export class ToolboxActivities { @Prop() graph: Graph; - //@State() activityCategoryModels: Array = []; private dnd: Addon.Dnd; - //private renderedActivities: Map; @State() private expandedCategories: Array = []; @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) { - // const browsableDescriptors = value.filter(x => x.isBrowsable); - // const categorizedActivitiesLookup = groupBy(browsableDescriptors, x => x.category); - // const categories = Object.keys(categorizedActivitiesLookup); - // const renderedActivities: Map = new Map(); - // - // // 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; diff --git a/src/modules/Elsa.Activities.Jobs/Activities/JobActivity.cs b/src/modules/Elsa.Activities.Jobs/Activities/JobActivity.cs index db0e4bc95..08318d4ee 100644 --- a/src/modules/Elsa.Activities.Jobs/Activities/JobActivity.cs +++ b/src/modules/Elsa.Activities.Jobs/Activities/JobActivity.cs @@ -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(); + 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); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Activities.Jobs/Attributes/JobAttribute.cs b/src/modules/Elsa.Activities.Jobs/Attributes/JobAttribute.cs new file mode 100644 index 000000000..4ebb6554e --- /dev/null +++ b/src/modules/Elsa.Activities.Jobs/Attributes/JobAttribute.cs @@ -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;} +} \ No newline at end of file diff --git a/src/modules/Elsa.Activities.Jobs/Elsa.Activities.Jobs.csproj b/src/modules/Elsa.Activities.Jobs/Elsa.Activities.Jobs.csproj index 41e3a7852..669c9b6b7 100644 --- a/src/modules/Elsa.Activities.Jobs/Elsa.Activities.Jobs.csproj +++ b/src/modules/Elsa.Activities.Jobs/Elsa.Activities.Jobs.csproj @@ -10,6 +10,7 @@ + diff --git a/src/modules/Elsa.Activities.Jobs/Handlers/JobExecutedHandler.cs b/src/modules/Elsa.Activities.Jobs/Handlers/JobExecutedHandler.cs new file mode 100644 index 000000000..25b191929 --- /dev/null +++ b/src/modules/Elsa.Activities.Jobs/Handlers/JobExecutedHandler.cs @@ -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 +{ + 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); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Activities.Jobs/Helpers/JobTypeNameHelper.cs b/src/modules/Elsa.Activities.Jobs/Helpers/JobTypeNameHelper.cs new file mode 100644 index 000000000..ba065eb6c --- /dev/null +++ b/src/modules/Elsa.Activities.Jobs/Helpers/JobTypeNameHelper.cs @@ -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(); + return activityAttr?.Namespace ?? jobType.Namespace; + } + + public static string GenerateTypeName(Type type, string? ns) + { + var activityAttr = type.GetCustomAttribute(); + var typeName = activityAttr?.TypeName ?? type.Name; + return ns != null ? $"{ns}.{typeName}" : typeName; + } + + public static string GenerateTypeName() 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)..]; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Activities.Jobs/Implementations/JobActivityProvider.cs b/src/modules/Elsa.Activities.Jobs/Implementations/JobActivityProvider.cs index ac3af7cff..066dab1bc 100644 --- a/src/modules/Elsa.Activities.Jobs/Implementations/JobActivityProvider.cs +++ b/src/modules/Elsa.Activities.Jobs/Implementations/JobActivityProvider.cs @@ -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> GetDescriptorsAsync(CancellationToken cancellationToken = default) + public ValueTask> GetDescriptorsAsync(CancellationToken cancellationToken = default) { var jobTypes = _jobRegistry.List(); var descriptors = CreateDescriptors(jobTypes).ToList(); - return descriptors; + return new(descriptors); } private IEnumerable CreateDescriptors(IEnumerable jobTypes) => jobTypes.Select(CreateDescriptor); private ActivityDescriptor CreateDescriptor(Type jobType) { - var typeName = jobType.Name; + var jobAttr = jobType.GetCustomAttribute(); + var ns = jobAttr?.Namespace ?? JobTypeNameHelper.GenerateNamespace(jobType); + var typeName = jobAttr?.TypeName ?? jobType.Name; + var fullTypeName = JobTypeNameHelper.GenerateTypeName(jobType); + var displayNameAttr = jobType.GetCustomAttribute(); + var displayName = displayNameAttr?.DisplayName ?? jobAttr?.DisplayName ?? typeName.Humanize(LetterCasing.Title); + var categoryAttr = jobType.GetCustomAttribute(); + var category = categoryAttr?.Category ?? jobAttr?.Category ?? ActivityTypeNameHelper.GetCategoryFromNamespace(ns) ?? "Miscellaneous"; + var descriptionAttr = jobType.GetCustomAttribute(); + 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(context); - activity.TypeName = typeName; + activity.TypeName = fullTypeName; activity.JobType = jobType; return activity; diff --git a/src/modules/Elsa.Activities.Jobs/Models/EnqueuedJobPayload.cs b/src/modules/Elsa.Activities.Jobs/Models/EnqueuedJobPayload.cs new file mode 100644 index 000000000..5d10a4b9a --- /dev/null +++ b/src/modules/Elsa.Activities.Jobs/Models/EnqueuedJobPayload.cs @@ -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!; +} \ No newline at end of file diff --git a/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs b/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs index 33356ff46..21905c0cc 100644 --- a/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs +++ b/src/modules/Elsa.AzureServiceBus/Activities/MessageReceived.cs @@ -55,7 +55,7 @@ public class MessageReceived : Trigger /// 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 if (!context.TryGetInput(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 context.Set(Result, body); } - private object GetBookmarkData(ExpressionExecutionContext context) + private object GetBookmarkPayload(ExpressionExecutionContext context) { var queueOrTopic = context.Get(QueueOrTopic)!; var subscription = context.Get(Subscription); diff --git a/src/modules/Elsa.AzureServiceBus/Handlers/UpdateWorkers.cs b/src/modules/Elsa.AzureServiceBus/Handlers/UpdateWorkers.cs index 874e4991d..3386c1d3c 100644 --- a/src/modules/Elsa.AzureServiceBus/Handlers/UpdateWorkers.cs +++ b/src/modules/Elsa.AzureServiceBus/Handlers/UpdateWorkers.cs @@ -14,10 +14,10 @@ namespace Elsa.AzureServiceBus.Handlers; /// public class UpdateWorkers : INotificationHandler, INotificationHandler, INotificationHandler { - 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; diff --git a/src/modules/Elsa.AzureServiceBus/Implementations/Worker.cs b/src/modules/Elsa.AzureServiceBus/Implementations/Worker.cs index b61c065f8..2afa49abc 100644 --- a/src/modules/Elsa.AzureServiceBus/Implementations/Worker.cs +++ b/src/modules/Elsa.AzureServiceBus/Implementations/Worker.cs @@ -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() { [MessageReceived.InputKey] = messageModel }; - var stimulus = Stimulus.Standard(BookmarkName, hash, input, correlationId); - var executionResults = (await _workflowService.DispatchStimulusAsync(stimulus, cancellationToken)).ToList(); + var input = new Dictionary { [MessageReceived.InputKey] = messageModel }; + var executionResults = (await _workflowService.DispatchStimulusAsync(BookmarkName, payload, input, correlationId, cancellationToken)).ToList(); _logger.LogInformation("Triggered {WorkflowCount} workflows", executionResults.Count); } diff --git a/src/modules/Elsa.Hangfire/Implementations/HangfireJobQueue.cs b/src/modules/Elsa.Hangfire/Implementations/HangfireJobQueue.cs index d5d6974ef..c8e85d27b 100644 --- a/src/modules/Elsa.Hangfire/Implementations/HangfireJobQueue.cs +++ b/src/modules/Elsa.Hangfire/Implementations/HangfireJobQueue.cs @@ -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 SubmitJobAsync(IJob job, string? queueName = default, CancellationToken cancellationToken = default) { var hangfireJob = HangfireJob.FromExpression(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); } } \ No newline at end of file diff --git a/src/modules/Elsa.Hangfire/Jobs/RunElsaJob.cs b/src/modules/Elsa.Hangfire/Jobs/RunElsaJob.cs index 544d4af44..8ca85b084 100644 --- a/src/modules/Elsa.Hangfire/Jobs/RunElsaJob.cs +++ b/src/modules/Elsa.Hangfire/Jobs/RunElsaJob.cs @@ -1,4 +1,5 @@ using Elsa.Jobs.Services; +using Hangfire.Server; namespace Elsa.Hangfire.Jobs; diff --git a/src/modules/Elsa.Http/Activities/HttpEndpoint.cs b/src/modules/Elsa.Http/Activities/HttpEndpoint.cs index 41c837982..f2eef3b1b 100644 --- a/src/modules/Elsa.Http/Activities/HttpEndpoint.cs +++ b/src/modules/Elsa.Http/Activities/HttpEndpoint.cs @@ -41,7 +41,7 @@ public class HttpEndpoint : Trigger )] public Input Policy { get; set; } = new(default(string?)); - protected override IEnumerable GetTriggerData(TriggerIndexingContext context) => GetBookmarkData(context.ExpressionExecutionContext); + protected override IEnumerable GetTriggerPayload(TriggerIndexingContext context) => GetBookmarkPayload(context.ExpressionExecutionContext); protected override void Execute(ActivityExecutionContext context) { @@ -49,7 +49,7 @@ public class HttpEndpoint : Trigger if (!context.TryGetInput(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 context.Set(Result, request); } - private IEnumerable GetBookmarkData(ExpressionExecutionContext context) + private IEnumerable GetBookmarkPayload(ExpressionExecutionContext context) { // Generate bookmark data for path and selected methods. var path = context.Get(Path); diff --git a/src/modules/Elsa.Scheduling/Activities/StartAt.cs b/src/modules/Elsa.Scheduling/Activities/StartAt.cs index 0521e6eea..07a02b2bf 100644 --- a/src/modules/Elsa.Scheduling/Activities/StartAt.cs +++ b/src/modules/Elsa.Scheduling/Activities/StartAt.cs @@ -41,7 +41,7 @@ public class StartAt : Trigger [Input] public Input 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); diff --git a/src/modules/Elsa.Scheduling/Activities/Timer.cs b/src/modules/Elsa.Scheduling/Activities/Timer.cs index 66f55cc9e..34ad1fbd2 100644 --- a/src/modules/Elsa.Scheduling/Activities/Timer.cs +++ b/src/modules/Elsa.Scheduling/Activities/Timer.cs @@ -29,7 +29,7 @@ public class Timer : EventGenerator [Input] public Input 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(); diff --git a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs index c704c97d3..beb46aeaa 100644 --- a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs +++ b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs @@ -65,7 +65,7 @@ public class WorkflowsFeature : FeatureBase .AddSingleton() .AddSingleton() .AddSingleton() - .AddSingleton() + .AddSingleton() .AddTransient() .AddSingleton(typeof(Func), sp => () => sp.GetRequiredService()) .AddSingleton() diff --git a/src/modules/Elsa.Workflows.Core/Implementations/BookmarkDataSerializer.cs b/src/modules/Elsa.Workflows.Core/Implementations/BookmarkPayloadSerializer.cs similarity index 85% rename from src/modules/Elsa.Workflows.Core/Implementations/BookmarkDataSerializer.cs rename to src/modules/Elsa.Workflows.Core/Implementations/BookmarkPayloadSerializer.cs index 9d1102e62..b303ffb3d 100644 --- a/src/modules/Elsa.Workflows.Core/Implementations/BookmarkDataSerializer.cs +++ b/src/modules/Elsa.Workflows.Core/Implementations/BookmarkPayloadSerializer.cs @@ -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 { diff --git a/src/modules/Elsa.Workflows.Core/Implementations/Hasher.cs b/src/modules/Elsa.Workflows.Core/Implementations/Hasher.cs index 30f7391ee..c7de4514f 100644 --- a/src/modules/Elsa.Workflows.Core/Implementations/Hasher.cs +++ b/src/modules/Elsa.Workflows.Core/Implementations/Hasher.cs @@ -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); } diff --git a/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs index 3fe3525d0..1c745e1f7 100644 --- a/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs @@ -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(); var identityGenerator = GetRequiredService(); - var bookmarkDataSerializer = GetRequiredService(); - var bookmarkDatumJson = bookmarkDatum != null ? bookmarkDataSerializer.Serialize(bookmarkDatum) : default; - var hash = bookmarkDatumJson != null ? hasher.Hash(bookmarkDatumJson) : default; + var payloadSerializer = GetRequiredService(); + 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); diff --git a/src/modules/Elsa.Workflows.Core/Models/Trigger.cs b/src/modules/Elsa.Workflows.Core/Models/Trigger.cs index d6aca407b..1ebddec03 100644 --- a/src/modules/Elsa.Workflows.Core/Models/Trigger.cs +++ b/src/modules/Elsa.Workflows.Core/Models/Trigger.cs @@ -26,12 +26,12 @@ public abstract class Trigger : Activity, ITrigger /// /// Override this method to return trigger data. /// - protected virtual IEnumerable GetTriggerData(TriggerIndexingContext context) => new[] { GetTriggerDatum(context) }; + protected virtual IEnumerable GetTriggerData(TriggerIndexingContext context) => new[] { GetTriggerPayload(context) }; /// /// Override this method to return a trigger datum. /// - protected virtual object GetTriggerDatum(TriggerIndexingContext context) => new(); + protected virtual object GetTriggerPayload(TriggerIndexingContext context) => new(); } public abstract class Trigger : Activity, ITrigger @@ -55,14 +55,14 @@ public abstract class Trigger : Activity, ITrigger /// protected virtual ValueTask> GetTriggerDataAsync(TriggerIndexingContext context) { - var hashes = GetTriggerData(context); + var hashes = GetTriggerPayload(context); return ValueTask.FromResult(hashes); } /// /// Override this method to return trigger data. /// - protected virtual IEnumerable GetTriggerData(TriggerIndexingContext context) => new[] { GetTriggerDatum(context) }; + protected virtual IEnumerable GetTriggerPayload(TriggerIndexingContext context) => new[] { GetTriggerDatum(context) }; /// /// Override this method to return a trigger datum. diff --git a/src/modules/Elsa.Workflows.Core/Services/IBookmarkDataSerializer.cs b/src/modules/Elsa.Workflows.Core/Services/IBookmarkPayloadSerializer.cs similarity index 77% rename from src/modules/Elsa.Workflows.Core/Services/IBookmarkDataSerializer.cs rename to src/modules/Elsa.Workflows.Core/Services/IBookmarkPayloadSerializer.cs index cfca217e6..560f6651c 100644 --- a/src/modules/Elsa.Workflows.Core/Services/IBookmarkDataSerializer.cs +++ b/src/modules/Elsa.Workflows.Core/Services/IBookmarkPayloadSerializer.cs @@ -1,6 +1,6 @@ namespace Elsa.Workflows.Core.Services; -public interface IBookmarkDataSerializer +public interface IBookmarkPayloadSerializer { T Deserialize(string json) where T : notnull; string Serialize(T payload) where T : notnull; diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowService.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowService.cs index 946c38e30..6c1b546f1 100644 --- a/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowService.cs +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowService.cs @@ -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 ExecuteWorkflowAsync(string definitionId, VersionOptions versionOptions, IDictionary? 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> 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> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, IDictionary 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> 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> DispatchStimulusAsync(IStimulus stimulus, CancellationToken cancellationToken) { diff --git a/src/modules/Elsa.Workflows.Runtime/Models/Stimulus.cs b/src/modules/Elsa.Workflows.Runtime/Models/Stimulus.cs index 18ce42f55..6ab0cbf1b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Models/Stimulus.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/Stimulus.cs @@ -17,6 +17,6 @@ public static class Stimulus new(ActivityTypeNameHelper.GenerateTypeName(), default, input.ToDictionary(), correlationId); public static StandardStimulus Standard(string activityTypeName, string? hash = default, IDictionary? 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); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowService.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowService.cs index bd9f70634..d7debb8fb 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowService.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowService.cs @@ -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 DispatchWorkflowAsync(string instanceId, Bookmark bookmark, IDictionary? input = default, string? correlationId = default, CancellationToken cancellationToken = default); Task> ExecuteStimulusAsync(IStimulus stimulus, CancellationToken cancellationToken = default); Task> DispatchStimulusAsync(IStimulus stimulus, CancellationToken cancellationToken = default); + Task> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, object inputs, string? correlationId = default, CancellationToken cancellationToken = default); + Task> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, IDictionary inputs, string? correlationId = default, CancellationToken cancellationToken = default); + Task> DispatchStimulusAsync(string bookmarkName, object bookmarkPayload, string? correlationId = default, CancellationToken cancellationToken = default); } \ No newline at end of file