diff --git a/Elsa.sln b/Elsa.sln index 38764b133..5df40b8c4 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -80,6 +80,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Serialization", "src\c EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Scripting.Liquid", "src\scripting\Elsa.Scripting.Liquid\Elsa.Scripting.Liquid.csproj", "{A7473430-F228-4534-AEA5-AA8D80E09364}" EndProject +Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "modules", "modules", "{5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Modules.Quartz", "src\modules\Elsa.Modules.Quartz\Elsa.Modules.Quartz.csproj", "{7D5A49B4-9A9B-496E-803B-DEB85B2C3132}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -170,6 +174,10 @@ Global {A7473430-F228-4534-AEA5-AA8D80E09364}.Debug|Any CPU.Build.0 = Debug|Any CPU {A7473430-F228-4534-AEA5-AA8D80E09364}.Release|Any CPU.ActiveCfg = Release|Any CPU {A7473430-F228-4534-AEA5-AA8D80E09364}.Release|Any CPU.Build.0 = Release|Any CPU + {7D5A49B4-9A9B-496E-803B-DEB85B2C3132}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {7D5A49B4-9A9B-496E-803B-DEB85B2C3132}.Debug|Any CPU.Build.0 = Debug|Any CPU + {7D5A49B4-9A9B-496E-803B-DEB85B2C3132}.Release|Any CPU.ActiveCfg = Release|Any CPU + {7D5A49B4-9A9B-496E-803B-DEB85B2C3132}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(NestedProjects) = preSolution {155227F0-A33B-40AA-A4B4-06F813EB921B} = {61017E64-6D00-49CB-9E81-5002DC8F7D5F} @@ -205,5 +213,7 @@ Global {820016B7-01CD-4032-8239-6A4E53F1A352} = {9F8AE7FB-E5F9-4DCB-9CF8-0362B3D18DAA} {55F9B33D-5FAF-4355-97EF-445C93208E85} = {C6658DE0-2B2F-47F0-BB61-2CA66D435C09} {A7473430-F228-4534-AEA5-AA8D80E09364} = {2633B8B9-4AEC-4A54-8832-1C940B54E041} + {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} = {61017E64-6D00-49CB-9E81-5002DC8F7D5F} + {7D5A49B4-9A9B-496E-803B-DEB85B2C3132} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} EndGlobalSection EndGlobal diff --git a/src/activities/Elsa.Activities.Http/HttpTrigger.cs b/src/activities/Elsa.Activities.Http/HttpTrigger.cs index 117ca9a91..e190bf6ef 100644 --- a/src/activities/Elsa.Activities.Http/HttpTrigger.cs +++ b/src/activities/Elsa.Activities.Http/HttpTrigger.cs @@ -2,13 +2,12 @@ using System.Linq; using System.Net.Http; using Elsa.Attributes; -using Elsa.Contracts; using Elsa.Management.Models; using Elsa.Models; namespace Elsa.Activities.Http; -public class HttpTrigger : Trigger +public class HttpTrigger : TriggerActivity { [Input] public Input Path { get; set; } = default!; @@ -20,32 +19,14 @@ public class HttpTrigger : Trigger [Output] public Output? Result { get; set; } - protected override IEnumerable GetHashInputs(TriggerIndexingContext context) + protected override IEnumerable GetHashInputs(TriggerIndexingContext context) => GetHashInputs(context.ExpressionExecutionContext); + protected override void Execute(ActivityExecutionContext context) => context.SetBookmarks(GetHashInputs(context.ExpressionExecutionContext)); + + private IEnumerable GetHashInputs(ExpressionExecutionContext context) { - var path = context.ExpressionExecutionContext.Get(Path); - var methods = context.ExpressionExecutionContext.Get(SupportedMethods); + // Generate a bookmark hash for path and selected methods. + var path = context.Get(Path); + var methods = context.Get(SupportedMethods); return methods!.Select(x => (path!.ToLowerInvariant(), x.ToLowerInvariant())).Cast().ToArray(); } - - protected override void Execute(ActivityExecutionContext context) - { - var bookmarks = CreateBookmarks(context).ToList(); - context.SetBookmarks(bookmarks); - } - - private IEnumerable CreateBookmarks(ActivityExecutionContext context) - { - var path = context.Get(Path)!; - var methods = context.Get(SupportedMethods)!; - var hasher = context.GetRequiredService(); - var identityGenerator = context.GetRequiredService(); - - foreach (var method in methods) - { - var hashInput = (path.ToLowerInvariant(), method.ToLowerInvariant()); - var hash = hasher.Hash(hashInput); - var bookmarkId = identityGenerator.GenerateId(); - yield return new Bookmark(bookmarkId, NodeType, hash, Id, context.Id); - } - } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Contracts/IJob.cs b/src/activities/Elsa.Activities.Scheduling/Contracts/IJob.cs new file mode 100644 index 000000000..02b52e2d2 --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Contracts/IJob.cs @@ -0,0 +1,10 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Activities.Scheduling.Contracts; + +public interface IJob +{ + string JobId { get; } + Task ExecuteAsync(CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Contracts/IJobScheduler.cs b/src/activities/Elsa.Activities.Scheduling/Contracts/IJobScheduler.cs new file mode 100644 index 000000000..d9c35349e --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Contracts/IJobScheduler.cs @@ -0,0 +1,9 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Activities.Scheduling.Contracts; + +public interface IJobScheduler +{ + Task ScheduleAsync(IJob job, ISchedule schedule, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Contracts/ISchedule.cs b/src/activities/Elsa.Activities.Scheduling/Contracts/ISchedule.cs new file mode 100644 index 000000000..1b6d47b76 --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Contracts/ISchedule.cs @@ -0,0 +1,5 @@ +namespace Elsa.Activities.Scheduling.Contracts; + +public interface ISchedule +{ +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Elsa.Activities.Scheduling.csproj b/src/activities/Elsa.Activities.Scheduling/Elsa.Activities.Scheduling.csproj index 6e698355e..6e3ee0fe4 100644 --- a/src/activities/Elsa.Activities.Scheduling/Elsa.Activities.Scheduling.csproj +++ b/src/activities/Elsa.Activities.Scheduling/Elsa.Activities.Scheduling.csproj @@ -10,4 +10,8 @@ + + + + diff --git a/src/activities/Elsa.Activities.Scheduling/Jobs/ResumeWorkflowJob.cs b/src/activities/Elsa.Activities.Scheduling/Jobs/ResumeWorkflowJob.cs new file mode 100644 index 000000000..3023244fa --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Jobs/ResumeWorkflowJob.cs @@ -0,0 +1,27 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.Scheduling.Contracts; + +namespace Elsa.Activities.Scheduling.Jobs; + +public class ResumeWorkflowJob : IJob +{ + public ResumeWorkflowJob() + { + } + + public ResumeWorkflowJob(string workflowInstanceId, string activityId) + { + WorkflowInstanceId = workflowInstanceId; + ActivityId = activityId; + } + + public string JobId => $"workflow-instance:{WorkflowInstanceId}-{ActivityId}"; + public string WorkflowInstanceId { get; init; } = default!; + public string ActivityId { get; init; } = default!; + + public Task ExecuteAsync(CancellationToken cancellationToken) + { + throw new System.NotImplementedException(); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Jobs/RunWorkflowJob.cs b/src/activities/Elsa.Activities.Scheduling/Jobs/RunWorkflowJob.cs new file mode 100644 index 000000000..76e97ff17 --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Jobs/RunWorkflowJob.cs @@ -0,0 +1,27 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.Scheduling.Contracts; +using Elsa.Models; + +namespace Elsa.Activities.Scheduling.Jobs; + +public class RunWorkflowJob : IJob +{ + public RunWorkflowJob() + { + } + + public RunWorkflowJob(WorkflowIdentity workflowIdentity) + { + WorkflowIdentity = workflowIdentity; + } + + public string JobId => $"workflow:{WorkflowIdentity.DefinitionId}"; + public WorkflowIdentity WorkflowIdentity { get; init; } = default!; + + + public Task ExecuteAsync(CancellationToken cancellationToken) + { + throw new System.NotImplementedException(); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Schedules/CronSchedule.cs b/src/activities/Elsa.Activities.Scheduling/Schedules/CronSchedule.cs new file mode 100644 index 000000000..0191beace --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Schedules/CronSchedule.cs @@ -0,0 +1,8 @@ +using Elsa.Activities.Scheduling.Contracts; + +namespace Elsa.Activities.Scheduling.Schedules; + +public class CronSchedule : ISchedule +{ + public string CronExpression { get; set; } = default!; +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Schedules/RecurringSchedule.cs b/src/activities/Elsa.Activities.Scheduling/Schedules/RecurringSchedule.cs new file mode 100644 index 000000000..bf1af0b07 --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Schedules/RecurringSchedule.cs @@ -0,0 +1,10 @@ +using System; +using Elsa.Activities.Scheduling.Contracts; + +namespace Elsa.Activities.Scheduling.Schedules; + +public class RecurringSchedule : ISchedule +{ + public DateTime StartAt { get; set; } + public TimeSpan Interval { get; set; } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Schedules/SpecificInstantSchedule.cs b/src/activities/Elsa.Activities.Scheduling/Schedules/SpecificInstantSchedule.cs new file mode 100644 index 000000000..e24662d2d --- /dev/null +++ b/src/activities/Elsa.Activities.Scheduling/Schedules/SpecificInstantSchedule.cs @@ -0,0 +1,9 @@ +using System; +using Elsa.Activities.Scheduling.Contracts; + +namespace Elsa.Activities.Scheduling.Schedules; + +public class SpecificInstantSchedule : ISchedule +{ + public DateTime DateTime { get; set; } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Scheduling/Timer.cs b/src/activities/Elsa.Activities.Scheduling/Timer.cs index 63b49d38b..6aee3fa2f 100644 --- a/src/activities/Elsa.Activities.Scheduling/Timer.cs +++ b/src/activities/Elsa.Activities.Scheduling/Timer.cs @@ -1,5 +1,7 @@ using System; +using System.Collections.Generic; using Elsa.Attributes; +using Elsa.Contracts; using Elsa.Models; namespace Elsa.Activities.Scheduling; @@ -7,4 +9,12 @@ namespace Elsa.Activities.Scheduling; public class Timer : Trigger { [Input] public Input Interval { get; set; } = default!; + + protected override IEnumerable GetHashInputs(TriggerIndexingContext context) + { + var interval = context.ExpressionExecutionContext.Get(Interval); + var clock = context.ExpressionExecutionContext.GetRequiredService(); + var executeAt = clock.UtcNow.Add(interval); + return new object[] { executeAt, interval }; + } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Contracts/ITrigger.cs b/src/core/Elsa.Core/Contracts/ITrigger.cs index 60c5afe65..7d7390300 100644 --- a/src/core/Elsa.Core/Contracts/ITrigger.cs +++ b/src/core/Elsa.Core/Contracts/ITrigger.cs @@ -4,6 +4,5 @@ namespace Elsa.Contracts; public interface ITrigger : INode { - string TriggerType { get; set; } ValueTask> GetHashInputsAsync(TriggerIndexingContext context, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/core/Elsa.Core/Models/ActivityExecutionContext.cs b/src/core/Elsa.Core/Models/ActivityExecutionContext.cs index aa5f4c20e..f76b52844 100644 --- a/src/core/Elsa.Core/Models/ActivityExecutionContext.cs +++ b/src/core/Elsa.Core/Models/ActivityExecutionContext.cs @@ -47,18 +47,35 @@ public class ActivityExecutionContext ScheduleActivity(activity, completionCallback); } + public void SetBookmarks(IEnumerable hashInputs, IDictionary? data = default, ExecuteActivityDelegate? callback = default) + { + foreach (var hashInput in hashInputs) + SetBookmark(hashInput, data, callback); + } + public void SetBookmarks(IEnumerable bookmarks) => _bookmarks.AddRange(bookmarks); public void SetBookmark(Bookmark bookmark) => _bookmarks.Add(bookmark); - public void SetBookmark(string? hash, IDictionary? data = default, ExecuteActivityDelegate? callback = default) => + public void SetBookmark(object? hashInput, IDictionary? data = default, ExecuteActivityDelegate? callback = default) + { + var hasher = GetRequiredService(); + var hash = hashInput != null ? hasher.Hash(hashInput) : default; + SetBookmark(hash, data, callback); + } + + public void SetBookmark(string? hash, IDictionary? data = default, ExecuteActivityDelegate? callback = default) + { + var identityGenerator = GetRequiredService(); + SetBookmark(new Bookmark( - Guid.NewGuid().ToString(), + identityGenerator.GenerateId(), Activity.NodeType, hash, Activity.Id, Id, data ?? new Dictionary(), callback?.Method.Name)); + } public T? GetProperty(string key) => Properties.TryGetValue(key, out var value) ? (T?)value : default; public void SetProperty(string key, T value) => Properties[key] = value; diff --git a/src/core/Elsa.Core/Models/ExpressionExecutionContext.cs b/src/core/Elsa.Core/Models/ExpressionExecutionContext.cs index 9dea97b22..f19032830 100644 --- a/src/core/Elsa.Core/Models/ExpressionExecutionContext.cs +++ b/src/core/Elsa.Core/Models/ExpressionExecutionContext.cs @@ -1,9 +1,14 @@ +using Microsoft.Extensions.DependencyInjection; + namespace Elsa.Models; public class ExpressionExecutionContext { - public ExpressionExecutionContext(Register register, ExpressionExecutionContext? parentContext) + private readonly IServiceProvider _serviceProvider; + + public ExpressionExecutionContext(IServiceProvider serviceProvider, Register register, ExpressionExecutionContext? parentContext) { + _serviceProvider = serviceProvider; Register = register; ParentContext = parentContext; } @@ -31,5 +36,6 @@ public class ExpressionExecutionContext Set(output.LocationReference, convertedValue); } + public T GetRequiredService() where T : notnull => _serviceProvider.GetRequiredService(); private RegisterLocation? GetLocationInternal(RegisterLocationReference locationReference) => Register.TryGetLocation(locationReference.Id, out var location) ? location : ParentContext?.GetLocationInternal(locationReference); } \ No newline at end of file diff --git a/src/core/Elsa.Core/Models/Trigger.cs b/src/core/Elsa.Core/Models/Trigger.cs index a45be0a98..1d9223bfc 100644 --- a/src/core/Elsa.Core/Models/Trigger.cs +++ b/src/core/Elsa.Core/Models/Trigger.cs @@ -1,22 +1,16 @@ using Elsa.Contracts; +using Elsa.Helpers; namespace Elsa.Models; -public class Trigger : Activity, ITrigger +public class Trigger : ITrigger { - protected Trigger() - { - } + protected Trigger() => NodeType = TypeNameHelper.GenerateTypeName(GetType()); + protected Trigger(string triggerType) => NodeType = triggerType; - protected Trigger(string triggerType) : base(triggerType) - { - } - - public string TriggerType - { - get => NodeType; - set => NodeType = value; - } + public string Id { get; set; } = default!; + public string NodeType { get; set; } + public IDictionary Metadata { get; set; } = new Dictionary(); public virtual ValueTask> GetHashInputsAsync(TriggerIndexingContext context, CancellationToken cancellationToken = default) { diff --git a/src/core/Elsa.Core/Models/TriggerActivity.cs b/src/core/Elsa.Core/Models/TriggerActivity.cs new file mode 100644 index 000000000..7b4986037 --- /dev/null +++ b/src/core/Elsa.Core/Models/TriggerActivity.cs @@ -0,0 +1,28 @@ +using Elsa.Contracts; + +namespace Elsa.Models; + +public class TriggerActivity : Activity, ITrigger +{ + protected TriggerActivity() + { + } + + protected TriggerActivity(string triggerType) : base(triggerType) + { + } + + public string TriggerType + { + get => NodeType; + set => NodeType = value; + } + + public virtual ValueTask> GetHashInputsAsync(TriggerIndexingContext context, CancellationToken cancellationToken = default) + { + var hashes = GetHashInputs(context); + return ValueTask.FromResult(hashes); + } + + protected virtual IEnumerable GetHashInputs(TriggerIndexingContext context) => Enumerable.Empty(); +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/ActivityInvoker.cs b/src/core/Elsa.Core/Services/ActivityInvoker.cs index 4ee899e37..8bfeb2a7a 100644 --- a/src/core/Elsa.Core/Services/ActivityInvoker.cs +++ b/src/core/Elsa.Core/Services/ActivityInvoker.cs @@ -6,10 +6,12 @@ namespace Elsa.Services; public class ActivityInvoker : IActivityInvoker { private readonly IActivityExecutionPipeline _pipeline; + private readonly IServiceProvider _serviceProvider; - public ActivityInvoker(IActivityExecutionPipeline pipeline) + public ActivityInvoker(IActivityExecutionPipeline pipeline, IServiceProvider serviceProvider) { _pipeline = pipeline; + _serviceProvider = serviceProvider; } public async Task InvokeAsync( @@ -25,7 +27,7 @@ public class ActivityInvoker : IActivityInvoker // Setup an activity execution context. var register = new Register(); - var expressionExecutionContext = new ExpressionExecutionContext(register, parentActivityExecutionContext?.ExpressionExecutionContext); + var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, parentActivityExecutionContext?.ExpressionExecutionContext); var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, parentActivityExecutionContext, expressionExecutionContext, activity, cancellationToken); // Declare locations. diff --git a/src/core/Elsa.Core/Services/WorkflowStateSerializer.cs b/src/core/Elsa.Core/Services/WorkflowStateSerializer.cs index 59c72e156..18fbd956a 100644 --- a/src/core/Elsa.Core/Services/WorkflowStateSerializer.cs +++ b/src/core/Elsa.Core/Services/WorkflowStateSerializer.cs @@ -9,6 +9,13 @@ namespace Elsa.Services; public class WorkflowStateSerializer : IWorkflowStateSerializer { + private readonly IServiceProvider _serviceProvider; + + public WorkflowStateSerializer(IServiceProvider serviceProvider) + { + _serviceProvider = serviceProvider; + } + public WorkflowState ReadState(WorkflowExecutionContext workflowExecutionContext) { var state = new WorkflowState @@ -119,7 +126,7 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer { var activity = workflowExecutionContext.FindActivityById(activityExecutionContextState.ScheduledActivityId); var register = new Register(activityExecutionContextState.Register.Locations); - var expressionExecutionContext = new ExpressionExecutionContext(register, default); + var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, default); var properties = activityExecutionContextState.Properties; var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, default, expressionExecutionContext, activity, workflowExecutionContext.CancellationToken) { diff --git a/src/core/Elsa.Management/Services/TriggerDescriber.cs b/src/core/Elsa.Management/Services/TriggerDescriber.cs index 45ca408bb..6fbbe2835 100644 --- a/src/core/Elsa.Management/Services/TriggerDescriber.cs +++ b/src/core/Elsa.Management/Services/TriggerDescriber.cs @@ -45,9 +45,9 @@ public class TriggerDescriber : ITriggerDescriber InputProperties = DescribeInputProperties(inputProperties).ToList(), Constructor = context => { - var activity = _activityFactory.Create(triggerType, context); - activity.TriggerType = fullTypeName; - return activity; + var trigger = _activityFactory.Create(triggerType, context); + trigger.NodeType = fullTypeName; + return trigger; } }; diff --git a/src/modules/Elsa.Modules.Quartz/Contracts/IElsaJobSerializer.cs b/src/modules/Elsa.Modules.Quartz/Contracts/IElsaJobSerializer.cs new file mode 100644 index 000000000..04d496934 --- /dev/null +++ b/src/modules/Elsa.Modules.Quartz/Contracts/IElsaJobSerializer.cs @@ -0,0 +1,10 @@ +using Elsa.Activities.Scheduling.Contracts; +using IElsaJob = Elsa.Activities.Scheduling.Contracts.IJob; + +namespace Elsa.Modules.Quartz.Contracts; + +public interface IElsaJobSerializer +{ + string Serialize(IJob job); + T Deserialize(string json) where T : IElsaJob; +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.Quartz/Elsa.Modules.Quartz.csproj b/src/modules/Elsa.Modules.Quartz/Elsa.Modules.Quartz.csproj new file mode 100644 index 000000000..deb966734 --- /dev/null +++ b/src/modules/Elsa.Modules.Quartz/Elsa.Modules.Quartz.csproj @@ -0,0 +1,19 @@ + + + + net6.0 + enable + enable + + + + + + + + + + + + + diff --git a/src/modules/Elsa.Modules.Quartz/Extensions/ServiceCollectionExtensions.cs b/src/modules/Elsa.Modules.Quartz/Extensions/ServiceCollectionExtensions.cs new file mode 100644 index 000000000..dda944b74 --- /dev/null +++ b/src/modules/Elsa.Modules.Quartz/Extensions/ServiceCollectionExtensions.cs @@ -0,0 +1,74 @@ +using Elsa.Activities.Scheduling.Contracts; +using Elsa.Activities.Scheduling.Jobs; +using Elsa.Modules.Quartz.Contracts; +using Elsa.Modules.Quartz.Jobs; +using Elsa.Modules.Quartz.Services; +using Microsoft.Extensions.DependencyInjection; +using Quartz; +using IElsaJob = Elsa.Activities.Scheduling.Contracts.IJob; + +namespace Elsa.Modules.Quartz.Extensions; + +public static class ServiceCollectionExtensions +{ + /// + /// This will register both Quartz as well as Elsa-specific services and jobs. + /// If you prefer to register Quartz yourself, use + /// + public static IServiceCollection AddQuartzModule( + this IServiceCollection services, + Action? configureQuartzOptions = default, + Action? configureQuartz = default, + Action? configureQuartzHostedService = default) + { + if (configureQuartzOptions != null) + services.Configure(configureQuartzOptions); + + return services + .AddQuartz(configure => + { + ConfigureQuartz(configure, configureQuartz); + ConfigureQuartzModule(services, configure); + }) + .AddQuartzHostedService(options => ConfigureQuartzHostedService(options, configureQuartzHostedService)); + } + + /// + /// This will register Elsa-specific services and jobs, but will **not** register Quartz itself. To register Quartz, you need to do so yourself, or use to register & configure Quartz for Elsa. + /// + public static IServiceCollection ConfigureQuartzModule(this IServiceCollection services, IServiceCollectionQuartzConfigurator quartz) + { + services + .AddSingleton() + .AddSingleton(); + + quartz.AddElsaJobs(); + + return services; + } + + private static IServiceCollectionQuartzConfigurator AddElsaJobs(this IServiceCollectionQuartzConfigurator quartz) + { + quartz.AddJob(); + + return quartz; + } + + private static void ConfigureQuartzHostedService(QuartzHostedServiceOptions options, Action? configureQuartzHostedService) + { + options.WaitForJobsToComplete = true; + configureQuartzHostedService?.Invoke(options); + } + + private static IServiceCollectionQuartzConfigurator AddJob(this IServiceCollectionQuartzConfigurator quartz) where TJob : IElsaJob => + quartz.AddJob>(job => job.StoreDurably().WithIdentity(nameof(RunWorkflowJob))); + + private static void ConfigureQuartz(IServiceCollectionQuartzConfigurator quartz, Action? configureQuartz) + { + quartz.UseMicrosoftDependencyInjectionJobFactory(); + quartz.AddJob>(job => job.StoreDurably().WithIdentity(nameof(RunWorkflowJob))); + quartz.UseSimpleTypeLoader(); + quartz.UseInMemoryStore(); + configureQuartz?.Invoke(quartz); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.Quartz/Jobs/QuartzJob.cs b/src/modules/Elsa.Modules.Quartz/Jobs/QuartzJob.cs new file mode 100644 index 000000000..4676288df --- /dev/null +++ b/src/modules/Elsa.Modules.Quartz/Jobs/QuartzJob.cs @@ -0,0 +1,25 @@ +using Elsa.Modules.Quartz.Contracts; +using Elsa.Modules.Quartz.Services; +using Quartz; +using IElsaJob = Elsa.Activities.Scheduling.Contracts.IJob; + +namespace Elsa.Modules.Quartz.Jobs; + +/// +/// A generic Quartz job that executes Elsa scheduled jobs. +/// +/// +public class QuartzJob : IJob where TElsaJob : IElsaJob +{ + private readonly IElsaJobSerializer _elsaJobSerializer; + + public QuartzJob(IElsaJobSerializer elsaJobSerializer) => _elsaJobSerializer = elsaJobSerializer; + + public async Task Execute(IJobExecutionContext context) + { + var json = context.MergedJobDataMap.GetString(QuartzJobScheduler.JobDataKey)!; + var elsaJob = _elsaJobSerializer.Deserialize(json); + + await elsaJob.ExecuteAsync(context.CancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.Quartz/Services/ElsaJobSerializer.cs b/src/modules/Elsa.Modules.Quartz/Services/ElsaJobSerializer.cs new file mode 100644 index 000000000..e2decbfac --- /dev/null +++ b/src/modules/Elsa.Modules.Quartz/Services/ElsaJobSerializer.cs @@ -0,0 +1,23 @@ +using System.Text.Json; +using Elsa.Activities.Scheduling.Contracts; +using Elsa.Modules.Quartz.Contracts; +using IElsaJob = Elsa.Activities.Scheduling.Contracts.IJob; + +namespace Elsa.Modules.Quartz.Services; + +public class ElsaJobSerializer : IElsaJobSerializer +{ + public string Serialize(IJob job) + { + var serializerOptions = CreateSerializerOptions(); + return JsonSerializer.Serialize(job, serializerOptions); + } + + public T Deserialize(string json) where T : IElsaJob + { + var serializerOptions = CreateSerializerOptions(); + return JsonSerializer.Deserialize(json, serializerOptions)!; + } + + private JsonSerializerOptions CreateSerializerOptions() => new(); +} \ No newline at end of file diff --git a/src/modules/Elsa.Modules.Quartz/Services/QuartzJobScheduler.cs b/src/modules/Elsa.Modules.Quartz/Services/QuartzJobScheduler.cs new file mode 100644 index 000000000..b1e700713 --- /dev/null +++ b/src/modules/Elsa.Modules.Quartz/Services/QuartzJobScheduler.cs @@ -0,0 +1,77 @@ +using Elsa.Activities.Scheduling.Schedules; +using Elsa.Modules.Quartz.Contracts; +using Elsa.Modules.Quartz.Jobs; +using Microsoft.Extensions.Logging; +using Quartz; +using IElsaJobScheduler = Elsa.Activities.Scheduling.Contracts.IJobScheduler; +using IElsaJob = Elsa.Activities.Scheduling.Contracts.IJob; +using IElsaSchedule = Elsa.Activities.Scheduling.Contracts.ISchedule; + +namespace Elsa.Modules.Quartz.Services; + +public class QuartzJobScheduler : IElsaJobScheduler +{ + public const string JobDataKey = "ElsaJob"; + private readonly IElsaJobSerializer _elsaJobSerializer; + private readonly ISchedulerFactory _schedulerFactory; + private readonly ILogger _logger; + + public QuartzJobScheduler(IElsaJobSerializer elsaJobSerializer, ISchedulerFactory schedulerFactory, ILogger logger) + { + _elsaJobSerializer = elsaJobSerializer; + _schedulerFactory = schedulerFactory; + _logger = logger; + } + + public async Task ScheduleAsync(IElsaJob job, IElsaSchedule schedule, CancellationToken cancellationToken = default) + { + var quartzTrigger = CreateTrigger(job, schedule); + await ScheduleJob(quartzTrigger, cancellationToken); + } + + private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) + { + var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); + + try + { + await scheduler.ScheduleJob(trigger, cancellationToken); + } + catch (SchedulerException e) + { + _logger.LogWarning(e, "Failed to schedule trigger {TriggerKey}", trigger.Key.ToString()); + } + } + + private ITrigger CreateTrigger(IElsaJob job, IElsaSchedule schedule) + { + var jobName = job.GetType().Name; + var json = _elsaJobSerializer.Serialize(job); + var builder = TriggerBuilder.Create().ForJob(jobName).WithIdentity(job.JobId).UsingJobData(JobDataKey, json); + + switch (schedule) + { + case RecurringSchedule recurringSchedule: + { + builder.StartAt(recurringSchedule.StartAt); + builder.WithSimpleSchedule(x => x.WithInterval(recurringSchedule.Interval).RepeatForever()); + break; + } + case CronSchedule cronSchedule: + { + builder.WithCronSchedule(cronSchedule.CronExpression); + break; + } + case SpecificInstantSchedule specificInstantSchedule: + { + builder.StartAt(specificInstantSchedule.DateTime); + break; + } + + default: + throw new NotSupportedException($"Schedule of type {schedule.GetType()} is not supported. But if you create an issue, we'll make this logic extensible & replaceable :)"); + } + + return builder.Build(); + } +} \ No newline at end of file diff --git a/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs b/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs index e7b319c36..716c1c479 100644 --- a/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs +++ b/src/runtime/Elsa.Runtime/Services/TriggerIndexer.cs @@ -17,6 +17,7 @@ public class TriggerIndexer : ITriggerIndexer private readonly IWorkflowRegistry _workflowRegistry; private readonly IExpressionEvaluator _expressionEvaluator; private readonly ICommandSender _mediator; + private readonly IServiceProvider _serviceProvider; private readonly IHasher _hasher; private readonly ILogger _logger; @@ -24,12 +25,14 @@ public class TriggerIndexer : ITriggerIndexer IWorkflowRegistry workflowRegistry, IExpressionEvaluator expressionEvaluator, ICommandSender mediator, + IServiceProvider serviceProvider, IHasher hasher, ILogger logger) { _workflowRegistry = workflowRegistry; _expressionEvaluator = expressionEvaluator; _mediator = mediator; + _serviceProvider = serviceProvider; _hasher = hasher; _logger = logger; } @@ -79,7 +82,7 @@ public class TriggerIndexer : ITriggerIndexer var inputs = trigger.GetInputs(); var assignedInputs = inputs.Where(x => x.LocationReference != null!).ToList(); var register = context.GetOrCreateRegister(trigger); - var expressionExecutionContext = new ExpressionExecutionContext(register, default); + var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, default); // Evaluate trigger inputs. foreach (var input in assignedInputs) @@ -121,9 +124,9 @@ public class TriggerIndexer : ITriggerIndexer } catch (Exception e) { - _logger.LogWarning( e, "Failed to get hash inputs"); + _logger.LogWarning(e, "Failed to get hash inputs"); } - return Array.Empty() ; + return Array.Empty(); } } \ No newline at end of file diff --git a/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj b/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj index 4684e150c..bf7c202be 100644 --- a/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj +++ b/src/samples/aspnet/Elsa.Samples.Web1/Elsa.Samples.Web1.csproj @@ -7,8 +7,10 @@ + + diff --git a/src/samples/aspnet/Elsa.Samples.Web1/Program.cs b/src/samples/aspnet/Elsa.Samples.Web1/Program.cs index a8523af01..0f394acbc 100644 --- a/src/samples/aspnet/Elsa.Samples.Web1/Program.cs +++ b/src/samples/aspnet/Elsa.Samples.Web1/Program.cs @@ -2,12 +2,14 @@ using Elsa.Activities.Console; using Elsa.Activities.ControlFlow; using Elsa.Activities.Http; using Elsa.Activities.Http.Extensions; +using Elsa.Activities.Scheduling; using Elsa.Activities.Workflows; using Elsa.Api.Extensions; using Elsa.Extensions; using Elsa.Management.Contracts; using Elsa.Management.Extensions; using Elsa.Mediator.Extensions; +using Elsa.Modules.Quartz.Extensions; using Elsa.Persistence.EntityFrameworkCore.Extensions; using Elsa.Persistence.EntityFrameworkCore.Sqlite; using Elsa.Persistence.Middleware.WorkflowExecution; @@ -52,17 +54,22 @@ services .AddActivity() .AddActivity() .AddActivity() - .AddActivity(); + .AddActivity() + ; // Register available triggers. services - .AddTrigger(); + .AddTrigger() + .AddTrigger(); // Register scripting languages. services .AddJavaScriptExpressions() .AddLiquidExpressions(); +// Register modules. +services.AddQuartzModule(); // Provides a scheduler implementation for Timer activities. + // Configure middleware pipeline. var app = builder.Build(); var serviceProvider = app.Services;