diff --git a/Elsa.sln b/Elsa.sln
index c2c52a26e..4711db729 100644
--- a/Elsa.sln
+++ b/Elsa.sln
@@ -170,6 +170,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Api.Client", "src\clie
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.QuartzIntegration", "src\samples\aspnet\Elsa.Samples.QuartzIntegration\Elsa.Samples.QuartzIntegration.csproj", "{B0312D9E-FA30-43E9-B666-40A8782D6E1C}"
EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.HangfireIntegration", "src\samples\aspnet\Elsa.Samples.HangfireIntegration\Elsa.Samples.HangfireIntegration.csproj", "{D2614FC7-102F-4F78-BB06-7C87304A10BA}"
+EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@@ -420,6 +422,10 @@ Global
{B0312D9E-FA30-43E9-B666-40A8782D6E1C}.Debug|Any CPU.Build.0 = Debug|Any CPU
{B0312D9E-FA30-43E9-B666-40A8782D6E1C}.Release|Any CPU.ActiveCfg = Release|Any CPU
{B0312D9E-FA30-43E9-B666-40A8782D6E1C}.Release|Any CPU.Build.0 = Release|Any CPU
+ {D2614FC7-102F-4F78-BB06-7C87304A10BA}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {D2614FC7-102F-4F78-BB06-7C87304A10BA}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {D2614FC7-102F-4F78-BB06-7C87304A10BA}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {D2614FC7-102F-4F78-BB06-7C87304A10BA}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(NestedProjects) = preSolution
{155227F0-A33B-40AA-A4B4-06F813EB921B} = {61017E64-6D00-49CB-9E81-5002DC8F7D5F}
@@ -494,5 +500,6 @@ Global
{89608AA5-5ADE-4832-AC7B-871C4AE64210} = {61017E64-6D00-49CB-9E81-5002DC8F7D5F}
{1FCB2200-28B8-4703-8E89-73241AAED047} = {89608AA5-5ADE-4832-AC7B-871C4AE64210}
{B0312D9E-FA30-43E9-B666-40A8782D6E1C} = {56C2FFB8-EA54-45B5-A095-4A78142EB4B5}
+ {D2614FC7-102F-4F78-BB06-7C87304A10BA} = {56C2FFB8-EA54-45B5-A095-4A78142EB4B5}
EndGlobalSection
EndGlobal
diff --git a/src/designer/designer_packages/elsa-workflows-designer/src/modules/notifications/models.ts b/src/designer/designer_packages/elsa-workflows-designer/src/modules/notifications/models.ts
index 6aaf74cf4..d87cc8c8e 100644
--- a/src/designer/designer_packages/elsa-workflows-designer/src/modules/notifications/models.ts
+++ b/src/designer/designer_packages/elsa-workflows-designer/src/modules/notifications/models.ts
@@ -3,7 +3,7 @@ import {Moment} from "moment";
export interface NotificationType {
id?: number | any;
title: string;
- text: string | JSX.Element;
+ text: string | any;
type?: NotificationDisplayType;
timestamp?: Moment;
showToast?: boolean;
diff --git a/src/modules/Elsa.Hangfire/Elsa.Hangfire.csproj b/src/modules/Elsa.Hangfire/Elsa.Hangfire.csproj
index b6d057e83..29acaa715 100644
--- a/src/modules/Elsa.Hangfire/Elsa.Hangfire.csproj
+++ b/src/modules/Elsa.Hangfire/Elsa.Hangfire.csproj
@@ -14,10 +14,12 @@
+
+
diff --git a/src/modules/Elsa.Hangfire/Extensions/JobStorageExtensions.cs b/src/modules/Elsa.Hangfire/Extensions/JobStorageExtensions.cs
new file mode 100644
index 000000000..4dce00947
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Extensions/JobStorageExtensions.cs
@@ -0,0 +1,32 @@
+using Hangfire;
+using Hangfire.Storage.Monitoring;
+
+namespace Elsa.Hangfire.Extensions;
+
+///
+/// A set of extension methods for .
+///
+public static class JobStorageExtensions
+{
+ ///
+ /// Enumerates all scheduled jobs of a given type.
+ ///
+ public static IEnumerable> EnumerateScheduledJobs(this JobStorage storage, string name)
+ {
+ var api = storage.GetMonitoringApi();
+ var skip = 0;
+ const int take = 100;
+ JobList jobList;
+
+ do
+ {
+ jobList = api.ScheduledJobs(skip, take);
+
+ var jobs = jobList.FindAll(x => x.Value.Job.Type == typeof(TJob));
+ foreach (var job in jobs.Where(x => (string)x.Value.Job.Args[0] == name))
+ yield return job;
+
+ skip += take;
+ } while (jobList.Count == take);
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Hangfire/Extensions/ModuleExtensions.cs
new file mode 100644
index 000000000..2b0efc3d1
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Extensions/ModuleExtensions.cs
@@ -0,0 +1,61 @@
+using Elsa.Features.Services;
+using Elsa.Hangfire.Features;
+using Elsa.Hangfire.Services;
+using Elsa.Scheduling.Contracts;
+using Elsa.Scheduling.Features;
+using Elsa.Workflows.Runtime.Features;
+using JetBrains.Annotations;
+
+// ReSharper disable once CheckNamespace
+namespace Elsa.Extensions;
+
+///
+/// Provides extension methods for the .
+///
+[PublicAPI]
+public static class ModuleExtensions
+{
+ ///
+ /// Installs and configures Hangfire. Only use this feature if you are not configuring Hangfire yourself.
+ ///
+ public static IModule UseHangfire(this IModule module, Action? configure = default)
+ {
+ return module.Use(configure);
+ }
+
+ ///
+ /// Configures Hangfire to use SQL Server storage. Only use this feature if you are not configuring Hangfire yourself.
+ ///
+ public static HangfireFeature UseSqlServerStorage(this HangfireFeature feature, Action configure)
+ {
+ feature.Module.Use(configure);
+ return feature;
+ }
+
+ ///
+ /// Configures Hangfire to use SQLite storage. Only use this feature if you are not configuring Hangfire yourself.
+ ///
+ public static HangfireFeature UseSqliteStorage(this HangfireFeature feature, Action configure)
+ {
+ feature.Module.Use(configure);
+ return feature;
+ }
+
+ ///
+ /// Installs a Hangfire implementation for .
+ ///
+ public static SchedulingFeature UseHangfireScheduler(this SchedulingFeature feature, Action? configure = default)
+ {
+ feature.Module.Use(configure);
+ return feature;
+ }
+
+ ///
+ /// Installs a Hangfire implementation for .
+ ///
+ public static WorkflowRuntimeFeature UseHangfireBackgroundActivityScheduler(this WorkflowRuntimeFeature feature, Action? configure = default)
+ {
+ feature.Module.Use(configure);
+ return feature;
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Extensions/TimeSpanExtensions.cs b/src/modules/Elsa.Hangfire/Extensions/TimeSpanExtensions.cs
new file mode 100644
index 000000000..9b2310038
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Extensions/TimeSpanExtensions.cs
@@ -0,0 +1,23 @@
+namespace Elsa.Hangfire.Extensions;
+
+///
+/// Adds extension methods to that converts it to a cron expression.
+///
+public static class TimeSpanExtensions
+{
+ ///
+ /// Converts the specified time span to a cron expression.
+ ///
+ /// The time span.
+ /// The cron expression.
+ public static string ToCronExpression(this TimeSpan timeSpan)
+ {
+ static string CreateCronComponent(int number) => (number > 0 ? $"*/{number}" : "*");
+
+ var cron = CreateCronComponent(timeSpan.Seconds);
+ cron += ' ' + CreateCronComponent(timeSpan.Minutes);
+ cron += ' ' + CreateCronComponent(timeSpan.Hours);
+ cron += ' ' + CreateCronComponent(timeSpan.Days);
+ return cron + " * *";
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Features/HangfireBackgroundActivityInvokerFeature.cs b/src/modules/Elsa.Hangfire/Features/HangfireBackgroundActivityInvokerFeature.cs
new file mode 100644
index 000000000..73ee284e9
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Features/HangfireBackgroundActivityInvokerFeature.cs
@@ -0,0 +1,38 @@
+using Elsa.Features.Abstractions;
+using Elsa.Features.Attributes;
+using Elsa.Features.Services;
+using Elsa.Hangfire.Services;
+using Elsa.Scheduling.Contracts;
+using Elsa.Scheduling.Features;
+using Elsa.Workflows.Runtime.Contracts;
+using Elsa.Workflows.Runtime.Features;
+using Microsoft.Extensions.DependencyInjection;
+
+namespace Elsa.Hangfire.Features;
+
+///
+/// Installs a Hangfire implementation for .
+///
+[DependsOn(typeof(WorkflowRuntimeFeature))]
+public class HangfireBackgroundActivitySchedulerFeature : FeatureBase
+{
+ ///
+ public HangfireBackgroundActivitySchedulerFeature(IModule module) : base(module)
+ {
+ }
+
+ ///
+ public override void Configure()
+ {
+ Module.Configure(workflowRuntimeFeature =>
+ {
+ workflowRuntimeFeature.BackgroundActivityInvoker = sp => sp.GetRequiredService();
+ });
+ }
+
+ ///
+ public override void Apply()
+ {
+ Services.AddSingleton();
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Features/HangfireFeature.cs b/src/modules/Elsa.Hangfire/Features/HangfireFeature.cs
index 04a8ada47..25c359cfe 100644
--- a/src/modules/Elsa.Hangfire/Features/HangfireFeature.cs
+++ b/src/modules/Elsa.Hangfire/Features/HangfireFeature.cs
@@ -2,13 +2,12 @@ using Elsa.Features.Abstractions;
using Elsa.Features.Services;
using Hangfire;
using Hangfire.MemoryStorage;
-using Hangfire.SqlServer;
using Newtonsoft.Json;
namespace Elsa.Hangfire.Features;
///
-/// Sets up Hangfire.
+/// Sets up Hangfire. If you're setting up Hangfire yourself, then you should not enable this feature.
///
public class HangfireFeature : FeatureBase
{
@@ -18,47 +17,39 @@ public class HangfireFeature : FeatureBase
}
///
- /// Whether to use SQL Server storage.
+ /// A delegate that configures Hangfire.
///
- public bool UseSqlServerStorage { get; set; }
+ public Action ConfigureHangfire { get; set; } = (_, cfg) => cfg.UseMemoryStorage();
+
+ ///
+ /// A delegate that configures Hangfire's background job server options.
+ ///
+ public Action ConfigureBackgroundServerOptions { get; set; } = (_, _) => { };
///
- /// The SQL Server storage options.
+ /// A delegate that creates a job storage instance.
///
- public SqlServerStorageOptions? SqlServerStorageOptions { get; set; }
-
- ///
- /// The SQL Server connection string.
- ///
- public string? SqlServerConnectionString { get; set; }
-
- ///
- /// The Hangfire background server options.
- ///
- public Action? ConfigureBackgroundServerOptions { get; set; }
+ public Func CreateJobStorage { get; set; } = () => new MemoryStorage();
///
- public override void Configure()
+ public override void Apply()
{
- Services.AddHangfire(configuration =>
+ Action configAction = (sp, cfg) =>
{
- configuration.UseSimpleAssemblyNameTypeSerializer();
- configuration.UseRecommendedSerializerSettings(json => json.TypeNameHandling = TypeNameHandling.Objects);
-
- if (UseSqlServerStorage)
- {
- var storageOptions = SqlServerStorageOptions ?? new SqlServerStorageOptions();
- configuration.UseSqlServerStorage(SqlServerConnectionString, storageOptions);
- }
- else
- {
- configuration.UseMemoryStorage();
- }
- });
-
- if (UseSqlServerStorage)
- Services.AddHangfireServer((_, options) => ConfigureBackgroundServerOptions?.Invoke(options), new SqlServerStorage(SqlServerConnectionString));
- else
- Services.AddHangfireServer(options => { ConfigureBackgroundServerOptions?.Invoke(options); });
+ cfg.UseSimpleAssemblyNameTypeSerializer();
+ cfg.UseRecommendedSerializerSettings(json => json.TypeNameHandling = TypeNameHandling.Objects);
+ };
+
+ Action serverOptionsAction = (sp, options) =>
+ {
+ options.WorkerCount = 1;
+ options.SchedulePollingInterval = TimeSpan.FromSeconds(1);
+ };
+
+ configAction += ConfigureHangfire;
+ serverOptionsAction += ConfigureBackgroundServerOptions;
+
+ Services.AddHangfire(configAction);
+ Services.AddHangfireServer(serverOptionsAction, CreateJobStorage());
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Features/HangfireSchedulerFeature.cs b/src/modules/Elsa.Hangfire/Features/HangfireSchedulerFeature.cs
new file mode 100644
index 000000000..c44223f0b
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Features/HangfireSchedulerFeature.cs
@@ -0,0 +1,33 @@
+using Elsa.Features.Abstractions;
+using Elsa.Features.Attributes;
+using Elsa.Features.Services;
+using Elsa.Hangfire.Services;
+using Elsa.Scheduling.Contracts;
+using Elsa.Scheduling.Features;
+using Microsoft.Extensions.DependencyInjection;
+
+namespace Elsa.Hangfire.Features;
+
+///
+/// Installs a Hangfire implementation for .
+///
+[DependsOn(typeof(SchedulingFeature))]
+public class HangfireSchedulerFeature : FeatureBase
+{
+ ///
+ public HangfireSchedulerFeature(IModule module) : base(module)
+ {
+ }
+
+ ///
+ public override void Configure()
+ {
+ Module.Configure(schedulingFeature => { schedulingFeature.WorkflowScheduler = sp => sp.GetRequiredService(); });
+ }
+
+ ///
+ public override void Apply()
+ {
+ Services.AddSingleton();
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Features/HangfireSqlServerStorageFeature.cs b/src/modules/Elsa.Hangfire/Features/HangfireSqlServerStorageFeature.cs
new file mode 100644
index 000000000..480a5a3bf
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Features/HangfireSqlServerStorageFeature.cs
@@ -0,0 +1,53 @@
+using Elsa.Extensions;
+using Elsa.Features.Abstractions;
+using Elsa.Features.Attributes;
+using Elsa.Features.Services;
+using Hangfire;
+using Hangfire.SqlServer;
+
+namespace Elsa.Hangfire.Features;
+
+///
+/// Configures the Hangfire feature to use SQL Server storage. If you're setting up Hangfire yourself, then you should not enable this feature.
+///
+[DependsOn(typeof(HangfireFeature))]
+public class HangfireSqlServerStorageFeature : FeatureBase
+{
+ ///
+ public HangfireSqlServerStorageFeature(IModule module) : base(module)
+ {
+ }
+
+ ///
+ /// The connection string to use when connecting to SQL Server, or the name of the connection string.
+ ///
+ public string NameOrConnectionString { get; set; } = default!;
+
+ ///
+ /// Configures the SQL Server storage options.
+ ///
+ public Action ConfigureSqlServerStorageOptions { get; set; } = _ => { };
+
+ ///
+ public override void Configure()
+ {
+ Module.Use(hangfireFeature =>
+ {
+ var storageOptions = new SqlServerStorageOptions
+ {
+ CommandBatchMaxTimeout = TimeSpan.FromMinutes(5),
+ SlidingInvisibilityTimeout = TimeSpan.FromMinutes(5),
+ QueuePollInterval = TimeSpan.FromSeconds(15),
+ UseRecommendedIsolationLevel = true
+ };
+ ConfigureSqlServerStorageOptions(storageOptions);
+
+ hangfireFeature.ConfigureHangfire = (_, cfg) =>
+ {
+ cfg.UseSqlServerStorage(NameOrConnectionString, storageOptions);
+ };
+
+ hangfireFeature.CreateJobStorage = () => new SqlServerStorage(NameOrConnectionString, storageOptions);
+ });
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Features/HangfireSqliteStorageFeature.cs b/src/modules/Elsa.Hangfire/Features/HangfireSqliteStorageFeature.cs
new file mode 100644
index 000000000..e039ac6e0
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Features/HangfireSqliteStorageFeature.cs
@@ -0,0 +1,51 @@
+using Elsa.Extensions;
+using Elsa.Features.Abstractions;
+using Elsa.Features.Attributes;
+using Elsa.Features.Services;
+using Hangfire;
+using Hangfire.SqlServer;
+using Hangfire.Storage.SQLite;
+
+namespace Elsa.Hangfire.Features;
+
+///
+/// Configures the Hangfire feature to use SQLite storage. If you're setting up Hangfire yourself, then you should not enable this feature.
+///
+[DependsOn(typeof(HangfireFeature))]
+public class HangfireSqliteStorageFeature : FeatureBase
+{
+ ///
+ public HangfireSqliteStorageFeature(IModule module) : base(module)
+ {
+ }
+
+ ///
+ /// The connection string to use when connecting to SQL Server, or the name of the connection string.
+ ///
+ public string NameOrConnectionString { get; set; } = default!;
+
+ ///
+ /// Configures the SQL Server storage options.
+ ///
+ public Action ConfigureSqlServerStorageOptions { get; set; } = _ => { };
+
+ ///
+ public override void Configure()
+ {
+ Module.Use(hangfireFeature =>
+ {
+ var storageOptions = new SQLiteStorageOptions
+ {
+ QueuePollInterval = TimeSpan.FromSeconds(1)
+ };
+ ConfigureSqlServerStorageOptions(storageOptions);
+
+ hangfireFeature.ConfigureHangfire = (_, cfg) =>
+ {
+ cfg.UseSQLiteStorage(NameOrConnectionString, storageOptions);
+ };
+
+ hangfireFeature.CreateJobStorage = () => new SQLiteStorage(NameOrConnectionString, storageOptions);
+ });
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Jobs/ExecuteBackgroundActivityJob.cs b/src/modules/Elsa.Hangfire/Jobs/ExecuteBackgroundActivityJob.cs
new file mode 100644
index 000000000..1793ce1e5
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Jobs/ExecuteBackgroundActivityJob.cs
@@ -0,0 +1,29 @@
+using Elsa.Workflows.Runtime.Contracts;
+using Elsa.Workflows.Runtime.Models;
+
+namespace Elsa.Hangfire.Jobs;
+
+///
+/// A job that executes a background activity.
+///
+public class ExecuteBackgroundActivityJob
+{
+ private readonly IBackgroundActivityInvoker _backgroundActivityInvoker;
+
+ ///
+ /// Initializes a new instance of the class.
+ ///
+ ///
+ public ExecuteBackgroundActivityJob(IBackgroundActivityInvoker backgroundActivityInvoker)
+ {
+ _backgroundActivityInvoker = backgroundActivityInvoker;
+ }
+
+ ///
+ /// Executes the job.
+ ///
+ public async Task ExecuteAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default)
+ {
+ await _backgroundActivityInvoker.ExecuteAsync(scheduledBackgroundActivity, cancellationToken);
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Jobs/ResumeWorkflowJob.cs b/src/modules/Elsa.Hangfire/Jobs/ResumeWorkflowJob.cs
new file mode 100644
index 000000000..76d4820f7
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Jobs/ResumeWorkflowJob.cs
@@ -0,0 +1,28 @@
+using Elsa.Workflows.Runtime.Contracts;
+using Elsa.Workflows.Runtime.Models.Requests;
+
+namespace Elsa.Hangfire.Jobs;
+
+///
+/// A job that resumes a workflow.
+///
+public class ResumeWorkflowJob
+{
+ private readonly IWorkflowDispatcher _workflowDispatcher;
+
+ ///
+ /// Initializes a new instance of the class.
+ ///
+ public ResumeWorkflowJob(IWorkflowDispatcher workflowDispatcher)
+ {
+ _workflowDispatcher = workflowDispatcher;
+ }
+
+ ///
+ /// Executes the job.
+ ///
+ /// The name of the job.
+ /// The workflow request.
+ /// The cancellation token.
+ public async Task ExecuteAsync(string name, DispatchWorkflowInstanceRequest request, CancellationToken cancellationToken) => await _workflowDispatcher.DispatchAsync(request, cancellationToken);
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Jobs/RunWorkflowJob.cs b/src/modules/Elsa.Hangfire/Jobs/RunWorkflowJob.cs
new file mode 100644
index 000000000..f8c20f943
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Jobs/RunWorkflowJob.cs
@@ -0,0 +1,28 @@
+using Elsa.Workflows.Runtime.Contracts;
+using Elsa.Workflows.Runtime.Models.Requests;
+
+namespace Elsa.Hangfire.Jobs;
+
+///
+/// A job that resumes a workflow.
+///
+public class RunWorkflowJob
+{
+ private readonly IWorkflowDispatcher _workflowDispatcher;
+
+ ///
+ /// Initializes a new instance of the class.
+ ///
+ public RunWorkflowJob(IWorkflowDispatcher workflowDispatcher)
+ {
+ _workflowDispatcher = workflowDispatcher;
+ }
+
+ ///
+ /// Executes the job.
+ ///
+ /// The name of the job.
+ /// The workflow request.
+ /// The cancellation token.
+ public async Task ExecuteAsync(string name, DispatchWorkflowDefinitionRequest request, CancellationToken cancellationToken) => await _workflowDispatcher.DispatchAsync(request, cancellationToken);
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Services/HangfireBackgroundActivityScheduler.cs b/src/modules/Elsa.Hangfire/Services/HangfireBackgroundActivityScheduler.cs
new file mode 100644
index 000000000..fa921d7d7
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Services/HangfireBackgroundActivityScheduler.cs
@@ -0,0 +1,29 @@
+using Elsa.Hangfire.Jobs;
+using Elsa.Workflows.Runtime.Contracts;
+using Elsa.Workflows.Runtime.Models;
+using Hangfire;
+
+namespace Elsa.Hangfire.Services;
+
+///
+/// Invokes activities from a background worker within the context of its workflow instance using Hangfire.
+///
+public class HangfireBackgroundActivityScheduler : IBackgroundActivityScheduler
+{
+ private readonly IBackgroundJobClient _backgroundJobClient;
+
+ ///
+ /// Initializes a new instance of the class.
+ ///
+ public HangfireBackgroundActivityScheduler(IBackgroundJobClient backgroundJobClient)
+ {
+ _backgroundJobClient = backgroundJobClient;
+ }
+
+ ///
+ public Task ScheduleAsync(ScheduledBackgroundActivity scheduledBackgroundActivity, CancellationToken cancellationToken = default)
+ {
+ var jobId = _backgroundJobClient.Enqueue(x => x.ExecuteAsync(scheduledBackgroundActivity, CancellationToken.None));
+ return Task.FromResult(jobId);
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.Hangfire/Services/HangfireWorkflowScheduler.cs b/src/modules/Elsa.Hangfire/Services/HangfireWorkflowScheduler.cs
new file mode 100644
index 000000000..263a586fc
--- /dev/null
+++ b/src/modules/Elsa.Hangfire/Services/HangfireWorkflowScheduler.cs
@@ -0,0 +1,99 @@
+using Elsa.Hangfire.Extensions;
+using Elsa.Hangfire.Jobs;
+using Elsa.Scheduling.Contracts;
+using Elsa.Workflows.Runtime.Models.Requests;
+using Hangfire;
+using Hangfire.Storage;
+
+namespace Elsa.Hangfire.Services;
+
+///
+/// An implementation of that uses Hangfire.
+///
+public class HangfireWorkflowScheduler : IWorkflowScheduler
+{
+ private readonly IBackgroundJobClient _backgroundJobClient;
+ private readonly IRecurringJobManager _recurringJobManager;
+ private readonly JobStorage _jobStorage;
+
+ ///
+ /// Initializes a new instance of the class.
+ ///
+ public HangfireWorkflowScheduler(IBackgroundJobClient backgroundJobClient, IRecurringJobManager recurringJobManager, JobStorage jobStorage)
+ {
+ _backgroundJobClient = backgroundJobClient;
+ _recurringJobManager = recurringJobManager;
+ _jobStorage = jobStorage;
+ }
+
+ ///
+ public ValueTask ScheduleAtAsync(string taskName, DispatchWorkflowDefinitionRequest request, DateTimeOffset at, CancellationToken cancellationToken = default)
+ {
+ _backgroundJobClient.Schedule(job => job.ExecuteAsync(taskName, request, CancellationToken.None), at);
+ return ValueTask.CompletedTask;
+ }
+
+ ///
+ public ValueTask ScheduleAtAsync(string taskName, DispatchWorkflowInstanceRequest request, DateTimeOffset at, CancellationToken cancellationToken = default)
+ {
+ _backgroundJobClient.Schedule(job => job.ExecuteAsync(taskName, request, CancellationToken.None), at);
+ return ValueTask.CompletedTask;
+ }
+
+ ///
+ public async ValueTask ScheduleRecurringAsync(string taskName, DispatchWorkflowDefinitionRequest request, DateTimeOffset startAt, TimeSpan interval, CancellationToken cancellationToken = default)
+ {
+ await ScheduleCronAsync(taskName, request, interval.ToCronExpression(), cancellationToken);
+ }
+
+ ///
+ public async ValueTask ScheduleRecurringAsync(string taskName, DispatchWorkflowInstanceRequest request, DateTimeOffset startAt, TimeSpan interval, CancellationToken cancellationToken = default)
+ {
+ await ScheduleCronAsync(taskName, request, interval.ToCronExpression(), cancellationToken);
+ }
+
+ ///
+ public ValueTask ScheduleCronAsync(string taskName, DispatchWorkflowDefinitionRequest request, string cronExpression, CancellationToken cancellationToken = default)
+ {
+ _recurringJobManager.AddOrUpdate(taskName, job => job.ExecuteAsync(taskName, request, CancellationToken.None), cronExpression);
+ return ValueTask.CompletedTask;
+ }
+
+ ///
+ public ValueTask ScheduleCronAsync(string taskName, DispatchWorkflowInstanceRequest request, string cronExpression, CancellationToken cancellationToken = default)
+ {
+ _recurringJobManager.AddOrUpdate(taskName, job => job.ExecuteAsync(taskName, request, CancellationToken.None), cronExpression);
+ return ValueTask.CompletedTask;
+ }
+
+ ///
+ public ValueTask UnscheduleAsync(string taskName, CancellationToken cancellationToken = default)
+ {
+ DeleteJobByTaskName(taskName);
+ return ValueTask.CompletedTask;
+ }
+
+ private void DeleteJobByTaskName(string taskName)
+ {
+ var scheduledJobIds = GetScheduledJobIds(taskName);
+ foreach (var jobId in scheduledJobIds) _backgroundJobClient.Delete(jobId);
+
+ var recurringJobIds = GetRecurringJobIds(taskName);
+ foreach (var jobId in recurringJobIds) _recurringJobManager.RemoveIfExists(jobId);
+ }
+
+ private IEnumerable GetScheduledJobIds(string taskName)
+ {
+ return _jobStorage.EnumerateScheduledJobs(taskName)
+ .Select(x => x.Key)
+ .Distinct()
+ .ToList();
+ }
+
+ private IEnumerable GetRecurringJobIds(string taskName)
+ {
+ using var connection = _jobStorage.GetConnection();
+ var jobs = connection.GetRecurringJobs().Where(x => x.Job.Type == typeof(TJob) && (string)x.Job.Args[0] == taskName);
+ return jobs.Select(x => x.Id).Distinct().ToList();
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs
index 7592d9d21..44ef0d16b 100644
--- a/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs
+++ b/src/modules/Elsa.MassTransit/Consumers/DispatchWorkflowRequestConsumer.cs
@@ -29,7 +29,7 @@ public class DispatchWorkflowRequestConsumer :
var message = context.Message;
var options = new StartWorkflowRuntimeOptions(message.CorrelationId, message.Input, message.VersionOptions, InstanceId: message.InstanceId);
- await _workflowRuntime.StartWorkflowAsync(message.DefinitionId, options, context.CancellationToken);
+ await _workflowRuntime.TryStartWorkflowAsync(message.DefinitionId, options, context.CancellationToken);
}
///
diff --git a/src/modules/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs b/src/modules/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs
index fc7a8a0a1..431a774c5 100644
--- a/src/modules/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs
+++ b/src/modules/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs
@@ -20,7 +20,7 @@ public class NotificationHandlerInvokerMiddleware : INotificationMiddleware
var notification = context.Notification;
var notificationType = notification.GetType();
var handlerType = typeof(INotificationHandler<>).MakeGenericType(notificationType);
- var handlers = _notificationHandlers.Where(x => handlerType.IsInstanceOfType(x)).ToArray();
+ var handlers = _notificationHandlers.Where(x => handlerType.IsInstanceOfType(x)).DistinctBy(x => x.GetType()).ToArray();
var handleMethod = handlerType.GetMethod("HandleAsync")!;
var cancellationToken = context.CancellationToken;
diff --git a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs
index 8ac84d7ef..379bffb6e 100644
--- a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs
+++ b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs
@@ -22,6 +22,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
private readonly ITriggerStore _triggerStore;
private readonly IIdentityGenerator _identityGenerator;
private readonly IBookmarkHasher _hasher;
+ private readonly IWorkflowDefinitionService _workflowDefinitionService;
private readonly IWorkflowInstanceFactory _workflowInstanceFactory;
///
@@ -33,6 +34,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
ITriggerStore triggerStore,
IIdentityGenerator identityGenerator,
IBookmarkHasher hasher,
+ IWorkflowDefinitionService workflowDefinitionService,
IWorkflowInstanceFactory workflowInstanceFactory)
{
_cluster = cluster;
@@ -40,6 +42,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
_triggerStore = triggerStore;
_identityGenerator = identityGenerator;
_hasher = hasher;
+ _workflowDefinitionService = workflowDefinitionService;
_workflowInstanceFactory = workflowInstanceFactory;
}
@@ -67,6 +70,18 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
return new CanStartWorkflowResult(workflowInstanceId, response!.CanStart);
}
+ ///
+ public async Task TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
+ {
+ // Load the workflow definition.
+ var workflowDefinition = await _workflowDefinitionService.FindAsync(definitionId, options.VersionOptions, cancellationToken);
+
+ if (workflowDefinition == null)
+ return null;
+
+ return await StartWorkflowAsync(definitionId, options, cancellationToken);
+ }
+
///
public async Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
{
@@ -177,7 +192,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
var collectedResumableWorkflow = (match as ResumableWorkflowMatch)!;
var runtimeOptions = new ResumeWorkflowRuntimeOptions(collectedResumableWorkflow.CorrelationId, Input: input);
-
+
var resumeResult = await ResumeWorkflowAsync(
match.WorkflowInstanceId,
runtimeOptions with { BookmarkId = collectedResumableWorkflow.BookmarkId },
diff --git a/src/modules/Elsa.Quartz/Extensions/JobDataMapExtensions.cs b/src/modules/Elsa.Quartz/Extensions/JobDataMapExtensions.cs
index fca4e77fd..983163a81 100644
--- a/src/modules/Elsa.Quartz/Extensions/JobDataMapExtensions.cs
+++ b/src/modules/Elsa.Quartz/Extensions/JobDataMapExtensions.cs
@@ -33,4 +33,13 @@ public static class JobDataMapExtensions
return map;
}
+
+ ///
+ /// Gets a dictionary from the map.
+ ///
+ public static IDictionary? GetDictionary(this JobDataMap map, string key)
+ {
+ var json = (string?)map.Get(key);
+ return json == null ? null : JsonSerializer.Deserialize>(json);
+ }
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Quartz/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Quartz/Extensions/ModuleExtensions.cs
index dd2789f64..3df03281e 100644
--- a/src/modules/Elsa.Quartz/Extensions/ModuleExtensions.cs
+++ b/src/modules/Elsa.Quartz/Extensions/ModuleExtensions.cs
@@ -7,7 +7,7 @@ using Elsa.Scheduling.Features;
namespace Elsa.Extensions;
///
-/// Provides extension methods for .
+/// Provides extension methods for the .
///
public static class ModuleExtensions
{
diff --git a/src/modules/Elsa.Quartz/Jobs/ResumeWorkflowTask.cs b/src/modules/Elsa.Quartz/Jobs/ResumeWorkflowJob.cs
similarity index 91%
rename from src/modules/Elsa.Quartz/Jobs/ResumeWorkflowTask.cs
rename to src/modules/Elsa.Quartz/Jobs/ResumeWorkflowJob.cs
index 06c1032f1..45be25828 100644
--- a/src/modules/Elsa.Quartz/Jobs/ResumeWorkflowTask.cs
+++ b/src/modules/Elsa.Quartz/Jobs/ResumeWorkflowJob.cs
@@ -1,3 +1,4 @@
+using Elsa.Extensions;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Models.Requests;
using Quartz;
@@ -35,7 +36,8 @@ public class ResumeWorkflowJob : IJob
ActivityId = (string?)map.Get(nameof(DispatchWorkflowInstanceRequest.ActivityId)),
ActivityNodeId = (string?)map.Get(nameof(DispatchWorkflowInstanceRequest.ActivityNodeId)),
ActivityInstanceId = (string?)map.Get(nameof(DispatchWorkflowInstanceRequest.ActivityInstanceId)),
- CorrelationId = (string?)map.Get(nameof(DispatchWorkflowInstanceRequest.CorrelationId))
+ CorrelationId = (string?)map.Get(nameof(DispatchWorkflowInstanceRequest.CorrelationId)),
+ Input = map.GetDictionary(nameof(DispatchWorkflowInstanceRequest.Input))
};
await _workflowDispatcher.DispatchAsync(request, context.CancellationToken);
}
diff --git a/src/modules/Elsa.Quartz/Jobs/RunWorkflowJob.cs b/src/modules/Elsa.Quartz/Jobs/RunWorkflowJob.cs
index 13b9e8e70..2ace7b340 100644
--- a/src/modules/Elsa.Quartz/Jobs/RunWorkflowJob.cs
+++ b/src/modules/Elsa.Quartz/Jobs/RunWorkflowJob.cs
@@ -1,4 +1,5 @@
using Elsa.Common.Models;
+using Elsa.Extensions;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Models.Requests;
using Quartz;
@@ -40,7 +41,8 @@ public class RunWorkflowJob : IJob
VersionOptions = VersionOptions.FromString((string)map.Get(nameof(DispatchWorkflowDefinitionRequest.VersionOptions))),
TriggerActivityId = (string?)map.Get(nameof(DispatchWorkflowDefinitionRequest.TriggerActivityId)),
InstanceId = (string?)map.Get(nameof(DispatchWorkflowDefinitionRequest.InstanceId)),
- CorrelationId = (string?)map.Get(nameof(DispatchWorkflowDefinitionRequest.CorrelationId))
+ CorrelationId = (string?)map.Get(nameof(DispatchWorkflowDefinitionRequest.CorrelationId)),
+ Input = map.GetDictionary(nameof(DispatchWorkflowDefinitionRequest.Input))
};
await _workflowDispatcher.DispatchAsync(request, context.CancellationToken);
}
diff --git a/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs
index 67b3c1b77..cfdadf4a3 100644
--- a/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs
+++ b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs
@@ -33,6 +33,7 @@ public class SchedulingFeature : FeatureBase
.AddSingleton()
.AddSingleton()
.AddSingleton()
+ .AddSingleton()
.AddSingleton(WorkflowScheduler)
.AddHandlersFrom();
diff --git a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs
index 094558141..d5e2e40e6 100644
--- a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs
+++ b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs
@@ -48,28 +48,36 @@ public class ScheduledRecurringTask : IScheduledTask
private void Schedule()
{
+ var startAt = _startAt;
+ var adjusted = false;
+
while (true)
{
var now = _systemClock.UtcNow;
- var delay = _startAt - now;
+ var delay = startAt - now;
- if (delay.Milliseconds <= 0)
+ if (!adjusted && delay.Milliseconds <= 0)
+ {
+ adjusted = true;
continue;
+ }
- SetupTimer(delay, now);
+ SetupTimer(delay);
break;
}
}
- private void SetupTimer(TimeSpan delay, DateTimeOffset now)
+ private void SetupTimer(TimeSpan delay)
{
+ if(delay < TimeSpan.Zero) delay = TimeSpan.FromSeconds(1);
+
_timer = new Timer(delay.TotalMilliseconds) { Enabled = true };
_timer.Elapsed += async (_, _) =>
{
_timer.Dispose();
_timer = null;
- _startAt = now + _interval;
+ _startAt = _systemClock.UtcNow + _interval;
var cancellationToken = _cancellationTokenSource.Token;
if (!cancellationToken.IsCancellationRequested) await _commandSender.SendAsync(new RunScheduledTask(_task), cancellationToken);
diff --git a/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs b/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs
index 859a43349..ee0ca79c7 100644
--- a/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs
+++ b/src/modules/Elsa.Scheduling/Services/DefaultTriggerScheduler.cs
@@ -77,7 +77,7 @@ public class DefaultTriggerScheduler : ITriggerScheduler
var timerTriggers = triggerList.Filter().ToList();
// Select all StartAt triggers.
- var startAtTriggers = triggerList.Filter().ToList();
+ var startAtTriggers = triggerList.Filter().ToList();
// Concatenate the filtered triggers.
var filteredTriggers = timerTriggers.Concat(startAtTriggers).ToList();
diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs
index a361766db..bcb406810 100644
--- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs
@@ -34,6 +34,11 @@ public interface IWorkflowRuntime
TriggerWorkflowsRuntimeOptions options,
CancellationToken cancellationToken = default);
+ ///
+ /// Tries to start a workflow and returns the result if successful.
+ ///
+ Task TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
+
///
/// Resumes an existing workflow instance.
///
diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs
index 352d91fc3..5aeac82c2 100644
--- a/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Handlers/DispatchWorkflowRequestHandler.cs
@@ -31,7 +31,7 @@ internal class DispatchWorkflowRequestHandler :
{
var options = new StartWorkflowRuntimeOptions(command.CorrelationId, command.Input, command.VersionOptions, InstanceId: command.InstanceId, TriggerActivityId: command.TriggerActivityId);
- await _workflowRuntime.StartWorkflowAsync(command.DefinitionId, options, cancellationToken);
+ await _workflowRuntime.TryStartWorkflowAsync(command.DefinitionId, options, cancellationToken);
return Unit.Instance;
}
diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs
index b3c98b981..1cfbe1947 100644
--- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs
+++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs
@@ -2,9 +2,11 @@ using Elsa.Common.Models;
using Elsa.Workflows.Core.Contracts;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.State;
+using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Models;
using Medallion.Threading;
+using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Runtime.Services;
@@ -21,6 +23,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
private readonly IBookmarkHasher _hasher;
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly IWorkflowInstanceFactory _workflowInstanceFactory;
+ private readonly ILogger _logger;
///
/// Constructor.
@@ -33,7 +36,8 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
IBookmarkStore bookmarkStore,
IBookmarkHasher hasher,
IDistributedLockProvider distributedLockProvider,
- IWorkflowInstanceFactory workflowInstanceFactory)
+ IWorkflowInstanceFactory workflowInstanceFactory,
+ ILogger logger)
{
_workflowHostFactory = workflowHostFactory;
_workflowDefinitionService = workflowDefinitionService;
@@ -43,6 +47,7 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
_hasher = hasher;
_distributedLockProvider = distributedLockProvider;
_workflowInstanceFactory = workflowInstanceFactory;
+ _logger = logger;
}
///
@@ -59,16 +64,19 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
///
public async Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
{
- var input = options.Input;
- var correlationId = options.CorrelationId;
var workflowHost = await CreateWorkflowHostAsync(definitionId, options, cancellationToken);
- var startWorkflowOptions = new StartWorkflowHostOptions(options.InstanceId, correlationId, input, options.TriggerActivityId);
- await workflowHost.StartWorkflowAsync(startWorkflowOptions, cancellationToken);
- var workflowState = workflowHost.WorkflowState;
+ return await StartWorkflowAsync(workflowHost, options, cancellationToken);
+ }
- await SaveWorkflowStateAsync(workflowState, cancellationToken);
+ ///
+ public async Task TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
+ {
+ var workflowDefinition = await FindWorkflowDefinitionAsync(definitionId, options.VersionOptions, cancellationToken);
- return new WorkflowExecutionResult(workflowState.Id, workflowState.Bookmarks);
+ if (workflowDefinition == null)
+ return null;
+
+ return await StartWorkflowAsync(workflowDefinition, options, cancellationToken);
}
///
@@ -125,7 +133,10 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
cancellationToken);
if (workflowDefinition == null)
- throw new Exception("Specified workflow definition and version does not exist");
+ {
+ _logger.LogInformation("The workflow definition {DefinitionId} version {Version} was not found", definitionId, version);
+ return new ResumeWorkflowResult(Array.Empty());
+ }
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var workflowHost = await _workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken);
@@ -215,14 +226,43 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
///
public async Task CountRunningWorkflowsAsync(CountRunningWorkflowsArgs args, CancellationToken cancellationToken = default) => await _workflowStateStore.CountAsync(args, cancellationToken);
+ private async Task StartWorkflowAsync(WorkflowDefinition workflowDefinition, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
+ {
+ var workflowHost = await CreateWorkflowHostAsync(workflowDefinition, cancellationToken);
+ return await StartWorkflowAsync(workflowHost, options, cancellationToken);
+ }
+
+ private async Task StartWorkflowAsync(IWorkflowHost workflowHost, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
+ {
+ var input = options.Input;
+ var correlationId = options.CorrelationId;
+ var startWorkflowOptions = new StartWorkflowHostOptions(options.InstanceId, correlationId, input, options.TriggerActivityId);
+ await workflowHost.StartWorkflowAsync(startWorkflowOptions, cancellationToken);
+ var workflowState = workflowHost.WorkflowState;
+
+ await SaveWorkflowStateAsync(workflowState, cancellationToken);
+
+ return new WorkflowExecutionResult(workflowState.Id, workflowState.Bookmarks);
+ }
+
+ private async Task FindWorkflowDefinitionAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken)
+ {
+ return await _workflowDefinitionService.FindAsync(definitionId, versionOptions, cancellationToken);
+ }
+
private async Task CreateWorkflowHostAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken)
{
var versionOptions = options.VersionOptions;
- var workflowDefinition = await _workflowDefinitionService.FindAsync(definitionId, versionOptions, cancellationToken);
+ var workflowDefinition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken);
if (workflowDefinition == null)
throw new Exception("Specified workflow definition and version does not exist");
-
+
+ return await CreateWorkflowHostAsync(workflowDefinition, cancellationToken);
+ }
+
+ private async Task CreateWorkflowHostAsync(WorkflowDefinition workflowDefinition, CancellationToken cancellationToken)
+ {
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
return await _workflowHostFactory.CreateAsync(workflow, cancellationToken);
}
diff --git a/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Elsa.Samples.HangfireIntegration.csproj b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Elsa.Samples.HangfireIntegration.csproj
new file mode 100644
index 000000000..0ff2cbf19
--- /dev/null
+++ b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Elsa.Samples.HangfireIntegration.csproj
@@ -0,0 +1,19 @@
+
+
+
+ net7.0
+ enable
+ enable
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Program.cs b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Program.cs
new file mode 100644
index 000000000..1ed50d5f7
--- /dev/null
+++ b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Program.cs
@@ -0,0 +1,69 @@
+using Elsa.EntityFrameworkCore.Extensions;
+using Elsa.EntityFrameworkCore.Modules.Management;
+using Elsa.EntityFrameworkCore.Modules.Runtime;
+using Elsa.Extensions;
+
+var builder = WebApplication.CreateBuilder(args);
+var configuration = builder.Configuration;
+var identitySection = configuration.GetSection("Identity");
+var identityTokenSection = identitySection.GetSection("Tokens");
+
+// Add Elsa to the container.
+builder.Services.AddElsa(elsa =>
+{
+ // Configure management feature to use EF Core.
+ elsa.UseWorkflowManagement(management => management.UseEntityFrameworkCore(ef => ef.UseSqlite()));
+
+ elsa.UseWorkflowRuntime(runtime =>
+ {
+ runtime.UseDefaultRuntime(dr => dr.UseEntityFrameworkCore(ef => ef.UseSqlite()));
+
+ // Use Hangfire to schedule background activities.
+ runtime.UseHangfireBackgroundActivityScheduler();
+
+ // Capture execution log records.
+ runtime.UseExecutionLogRecords(e => e.UseEntityFrameworkCore(ef => ef.UseSqlite()));
+
+ // Capture workflow state.
+ runtime.UseAsyncWorkflowStateExporter();
+ });
+
+ // Expose API endpoints.
+ elsa.UseWorkflowsApi();
+
+ // Use Hangfire.
+ elsa.UseHangfire(hangfire => hangfire.UseSqliteStorage(sqlite => sqlite.NameOrConnectionString = "elsa.sqlite.db"));
+
+ // Use hangfire for scheduling timer events.
+ elsa.UseScheduling(scheduling => scheduling.UseHangfireScheduler());
+
+ // Configure identity.
+ elsa.UseIdentity(identity =>
+ {
+ identity.IdentityOptions = options => identitySection.Bind(options);
+ identity.TokenOptions = options => identityTokenSection.Bind(options);
+ identity.UseConfigurationBasedUserProvider(options => identitySection.Bind(options));
+ identity.UseConfigurationBasedApplicationProvider(options => identitySection.Bind(options));
+ identity.UseConfigurationBasedRoleProvider(options => identitySection.Bind(options));
+ });
+
+ // Use default authentication (JWT).
+ elsa.UseDefaultAuthentication();
+});
+
+// Configure CORS to allow designer app hosted on a different origin to invoke the APIs.
+builder.Services.AddCors(cors => cors.AddDefaultPolicy(policy => policy.AllowAnyOrigin().AllowAnyHeader().AllowAnyMethod()));
+
+// Build the web app.
+var app = builder.Build();
+
+// Configure the web app's request pipeline.
+app.UseHttpsRedirection();
+app.UseCors();
+app.UseAuthentication();
+app.UseAuthorization();
+app.UseWorkflowsApi();
+app.UseWorkflows();
+
+// Run the web app.
+app.Run();
\ No newline at end of file
diff --git a/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Properties/launchSettings.json b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Properties/launchSettings.json
new file mode 100644
index 000000000..3513ccba1
--- /dev/null
+++ b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/Properties/launchSettings.json
@@ -0,0 +1,37 @@
+{
+ "iisSettings": {
+ "windowsAuthentication": false,
+ "anonymousAuthentication": true,
+ "iisExpress": {
+ "applicationUrl": "http://localhost:3978",
+ "sslPort": 44367
+ }
+ },
+ "profiles": {
+ "http": {
+ "commandName": "Project",
+ "dotnetRunMessages": true,
+ "launchBrowser": true,
+ "applicationUrl": "http://localhost:5090",
+ "environmentVariables": {
+ "ASPNETCORE_ENVIRONMENT": "Development"
+ }
+ },
+ "https": {
+ "commandName": "Project",
+ "dotnetRunMessages": true,
+ "launchBrowser": true,
+ "applicationUrl": "https://localhost:7020;http://localhost:5090",
+ "environmentVariables": {
+ "ASPNETCORE_ENVIRONMENT": "Development"
+ }
+ },
+ "IIS Express": {
+ "commandName": "IISExpress",
+ "launchBrowser": true,
+ "environmentVariables": {
+ "ASPNETCORE_ENVIRONMENT": "Development"
+ }
+ }
+ }
+}
diff --git a/src/samples/aspnet/Elsa.Samples.HangfireIntegration/README.md b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/README.md
new file mode 100644
index 000000000..94f9fae21
--- /dev/null
+++ b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/README.md
@@ -0,0 +1,6 @@
+## Secrets
+The following are the secrets stored in hashed form in appsettings.json:
+
+**API key**: `4E753976726458745954355043687772-e54d5a2c-33a3-4c05-a216-b09569062aed`
+**Admin user**: `admin`
+**Admin password**: `password`
\ No newline at end of file
diff --git a/src/samples/aspnet/Elsa.Samples.HangfireIntegration/appsettings.json b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/appsettings.json
new file mode 100644
index 000000000..9f2c93602
--- /dev/null
+++ b/src/samples/aspnet/Elsa.Samples.HangfireIntegration/appsettings.json
@@ -0,0 +1,44 @@
+{
+ "Logging": {
+ "LogLevel": {
+ "Default": "Warning",
+ "Microsoft.Hosting": "Information",
+ "Hangfire": "Warning"
+ }
+ },
+ "AllowedHosts": "*",
+ "Identity": {
+ "Tokens": {
+ "SigningKey": "secret-signing-key",
+ "AccessTokenLifetime": "1:00:00:00",
+ "RefreshTokenLifetime": "1:00:10:00"
+ },
+ "Roles": [{
+ "Id": "admin",
+ "Name": "Administrator",
+ "Permissions": ["*"]
+ }],
+ "Users": [
+ {
+ "Id": "a2323f46-42db-4e15-af8b-94238717d817",
+ "Name": "admin",
+ "HashedPassword": "TfKzh9RLix6FPcCNeHLkGrysFu3bYxqzGqduNdi8v1U=",
+ "HashedPasswordSalt": "JEy9kBlhHCNsencitRHlGxmErmSgY+FVyMJulCH27Ds=",
+ "Roles": ["admin"]
+ }
+ ],
+ "Applications": [{
+ "id": "529572c2df854b13807b8bf23f1784cd",
+ "name": "Postman",
+ "roles": [
+ "admin"
+ ],
+ "clientId": "Nu9vrdXtYT5PChwr",
+ "clientSecret": "011pp2C$|j01-qrMZpC9VC0F00XCJq(5",
+ "hashedApiKey": "d0rDld3A+ugKmdctGtMzOLTYjQFkOlUWN+kt0VyW9D0=",
+ "hashedApiKeySalt": "EnutGOyy5MuJWV0fF5jCQiciK7a8PU/DRF+fr6nekSY=",
+ "hashedClientSecret": "ERia2zBcCSWb/9dvB0grQ9yf7fWgFrClNeR8A5RMTzk=",
+ "hashedClientSecretSalt": "z3z8KmzHt+xkAj/zYTXcB8I7y0xAkLm95v4Er/oNqiY="
+ }]
+ }
+}