WIP
This commit is contained in:
parent
7814be98b3
commit
030634643c
|
|
@ -12,9 +12,9 @@ namespace Elsa.Activities.Temporal
|
|||
[Trigger(Category = "Timers", Description = "Cancel a timer (Cron, StartAt, Timer) so that it is not executed.")]
|
||||
public class ClearTimer : Activity
|
||||
{
|
||||
private readonly IWorkflowScheduler _workflowScheduler;
|
||||
private readonly IWorkflowInstanceScheduler _workflowScheduler;
|
||||
|
||||
public ClearTimer(IWorkflowScheduler workflowScheduler)
|
||||
public ClearTimer(IWorkflowInstanceScheduler workflowScheduler)
|
||||
{
|
||||
_workflowScheduler = workflowScheduler;
|
||||
}
|
||||
|
|
@ -24,7 +24,7 @@ namespace Elsa.Activities.Temporal
|
|||
|
||||
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context)
|
||||
{
|
||||
await _workflowScheduler.UnscheduleWorkflowAsync(null, context.WorkflowInstance.Id, ActivityId, context.WorkflowInstance.TenantId, context.CancellationToken);
|
||||
await _workflowScheduler.UnscheduleAsync(context.WorkflowInstance.Id, ActivityId, context.CancellationToken);
|
||||
return Done();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -20,10 +20,10 @@ namespace Elsa.Activities.Temporal
|
|||
{
|
||||
private readonly IClock _clock;
|
||||
private readonly IWorkflowInstanceStore _workflowInstanceStore;
|
||||
private readonly IWorkflowScheduler _workflowScheduler;
|
||||
private readonly IWorkflowInstanceScheduler _workflowScheduler;
|
||||
private readonly ICrontabParser _crontabParser;
|
||||
|
||||
public Cron(IWorkflowInstanceStore workflowInstanceStore, IWorkflowScheduler workflowScheduler, ICrontabParser crontabParser, IClock clock)
|
||||
public Cron(IWorkflowInstanceStore workflowInstanceStore, IWorkflowInstanceScheduler workflowScheduler, ICrontabParser crontabParser, IClock clock)
|
||||
{
|
||||
_clock = clock;
|
||||
_workflowInstanceStore = workflowInstanceStore;
|
||||
|
|
@ -56,7 +56,7 @@ namespace Elsa.Activities.Temporal
|
|||
return Done();
|
||||
|
||||
await _workflowInstanceStore.SaveAsync(context.WorkflowExecutionContext.WorkflowInstance, cancellationToken);
|
||||
await _workflowScheduler.ScheduleWorkflowAsync(null, workflowInstance.Id, Id, tenantId, executeAt, null, cancellationToken);
|
||||
await _workflowScheduler.ScheduleAsync(workflowInstance.Id, Id, executeAt, null, cancellationToken);
|
||||
|
||||
return Suspend();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
using Elsa.Activities.Temporal.Common.ActivityResults;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.ActivityResults;
|
||||
using Elsa.Attributes;
|
||||
using Elsa.Design;
|
||||
using Elsa.Expressions;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
|
|
@ -24,6 +26,14 @@ namespace Elsa.Activities.Temporal
|
|||
|
||||
[ActivityProperty(Hint = "The time interval at which this activity should tick.", SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })]
|
||||
public Duration Timeout { get; set; } = default!;
|
||||
|
||||
[ActivityProperty(
|
||||
Hint = "Whether workflows starting with this timer should execute on one node or on all nodes in a cluster.",
|
||||
SupportedSyntaxes = new[] { SyntaxNames.Literal, SyntaxNames.JavaScript, SyntaxNames.Liquid },
|
||||
Category = PropertyCategories.Hosting,
|
||||
DefaultValue = Common.Options.ClusterMode.SingleNode
|
||||
)]
|
||||
public ClusterMode ClusterMode { get; set; } = ClusterMode.SingleNode;
|
||||
|
||||
public Instant? ExecuteAt
|
||||
{
|
||||
|
|
|
|||
|
|
@ -26,8 +26,8 @@ namespace Elsa.Activities.Temporal.Common.ActivityResults
|
|||
|
||||
async ValueTask ScheduleWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken)
|
||||
{
|
||||
var scheduler = workflowExecutionContext.ServiceProvider.GetRequiredService<IWorkflowScheduler>();
|
||||
await scheduler.ScheduleWorkflowAsync(null, workflowInstanceId, activityId, tenantId, executeAt, null, cancellationToken);
|
||||
var scheduler = workflowExecutionContext.ServiceProvider.GetRequiredService<IWorkflowInstanceScheduler>();
|
||||
await scheduler.ScheduleAsync(workflowInstanceId, activityId, executeAt, null, cancellationToken);
|
||||
}
|
||||
|
||||
activityExecutionContext.WorkflowExecutionContext.RegisterTask(ScheduleWorkflowAsync);
|
||||
|
|
|
|||
|
|
@ -2,6 +2,8 @@
|
|||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.Attributes;
|
||||
using Elsa.Bookmarks;
|
||||
using NodaTime;
|
||||
|
||||
|
|
@ -11,6 +13,9 @@ namespace Elsa.Activities.Temporal.Common.Bookmarks
|
|||
{
|
||||
public Instant ExecuteAt { get; set; }
|
||||
public Duration Interval { get; set; }
|
||||
|
||||
[ExcludeFromHash]
|
||||
public ClusterMode? ClusterMode { get; set; }
|
||||
}
|
||||
|
||||
public class TimerBookmarkProvider : BookmarkProvider<TimerBookmark, Timer>
|
||||
|
|
@ -25,6 +30,7 @@ namespace Elsa.Activities.Temporal.Common.Bookmarks
|
|||
public override async ValueTask<IEnumerable<IBookmark>> GetBookmarksAsync(BookmarkProviderContext<Timer> context, CancellationToken cancellationToken)
|
||||
{
|
||||
var interval = await context.Activity.GetPropertyValueAsync(x => x.Timeout, cancellationToken);
|
||||
var clusterMode = context.Mode == BookmarkIndexingMode.WorkflowBlueprint ? await context.Activity.GetPropertyValueAsync(x => x.ClusterMode, cancellationToken) : default(ClusterMode?);
|
||||
var executeAt = GetExecuteAt(context, interval);
|
||||
|
||||
if (executeAt != null)
|
||||
|
|
@ -33,7 +39,8 @@ namespace Elsa.Activities.Temporal.Common.Bookmarks
|
|||
new TimerBookmark
|
||||
{
|
||||
ExecuteAt = executeAt.Value,
|
||||
Interval = interval
|
||||
Interval = interval,
|
||||
ClusterMode = clusterMode
|
||||
}
|
||||
};
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,23 @@
|
|||
using System.Threading.Tasks;
|
||||
using Elsa.Events;
|
||||
using MediatR;
|
||||
using Rebus.Handlers;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Common.Consumers
|
||||
{
|
||||
public class RemoveScheduledTriggersConsumer : IHandleMessages<WorkflowDefinitionPublished>, IHandleMessages<WorkflowDefinitionRetracted>, IHandleMessages<WorkflowDefinitionDeleted>
|
||||
{
|
||||
private readonly IMediator _mediator;
|
||||
|
||||
public RemoveScheduledTriggersConsumer(IMediator mediator)
|
||||
{
|
||||
_mediator = mediator;
|
||||
}
|
||||
|
||||
public Task Handle(WorkflowDefinitionPublished message) => _mediator.Publish(message);
|
||||
|
||||
public Task Handle(WorkflowDefinitionRetracted message) => _mediator.Publish(message);
|
||||
|
||||
public Task Handle(WorkflowDefinitionDeleted message) => _mediator.Publish(message);
|
||||
}
|
||||
}
|
||||
|
|
@ -25,4 +25,9 @@
|
|||
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<Folder Include="Models" />
|
||||
<Folder Include="Serialization" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
|
|||
|
|
@ -1,8 +1,10 @@
|
|||
using System;
|
||||
using Elsa.Activities.Temporal.Common.Bookmarks;
|
||||
using Elsa.Activities.Temporal.Common.Consumers;
|
||||
using Elsa.Activities.Temporal.Common.Handlers;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.Activities.Temporal.Common.StartupTasks;
|
||||
using Elsa.Events;
|
||||
using Elsa.Runtime;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
|
||||
|
|
@ -42,6 +44,10 @@ namespace Elsa.Activities.Temporal
|
|||
.AddActivity<StartAt>()
|
||||
.AddActivity<ClearTimer>();
|
||||
|
||||
options.AddPubSubConsumer<RemoveScheduledTriggersConsumer, WorkflowDefinitionPublished>();
|
||||
options.AddPubSubConsumer<RemoveScheduledTriggersConsumer, WorkflowDefinitionRetracted>();
|
||||
options.AddPubSubConsumer<RemoveScheduledTriggersConsumer, WorkflowDefinitionDeleted>();
|
||||
|
||||
return options;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -6,24 +6,29 @@ using MediatR;
|
|||
|
||||
namespace Elsa.Activities.Temporal.Common.Handlers
|
||||
{
|
||||
public class RemoveScheduledTriggers : INotificationHandler<BlockingActivityRemoved>, INotificationHandler<WorkflowDefinitionPublished>, INotificationHandler<WorkflowDefinitionRetracted>
|
||||
public class RemoveScheduledTriggers : INotificationHandler<BlockingActivityRemoved>, INotificationHandler<WorkflowDefinitionPublished>, INotificationHandler<WorkflowDefinitionRetracted>, INotificationHandler<WorkflowDefinitionDeleted>
|
||||
{
|
||||
private readonly IWorkflowScheduler _workflowScheduler;
|
||||
public RemoveScheduledTriggers(IWorkflowScheduler workflowScheduler) => _workflowScheduler = workflowScheduler;
|
||||
private readonly IWorkflowDefinitionScheduler _workflowDefinitionScheduler;
|
||||
private readonly IWorkflowInstanceScheduler _workflowInstanceScheduler;
|
||||
|
||||
public RemoveScheduledTriggers(IWorkflowDefinitionScheduler workflowDefinitionScheduler, IWorkflowInstanceScheduler workflowInstanceScheduler)
|
||||
{
|
||||
_workflowDefinitionScheduler = workflowDefinitionScheduler;
|
||||
_workflowInstanceScheduler = workflowInstanceScheduler;
|
||||
}
|
||||
|
||||
public async Task Handle(BlockingActivityRemoved notification, CancellationToken cancellationToken)
|
||||
{
|
||||
// TODO: Consider introducing a "stereotype" field for activities to exit early in case they are not stereotyped as "temporal".
|
||||
|
||||
await _workflowScheduler.UnscheduleWorkflowAsync(
|
||||
null,
|
||||
await _workflowInstanceScheduler.UnscheduleAsync(
|
||||
notification.WorkflowExecutionContext.WorkflowInstance.Id,
|
||||
notification.BlockingActivity.ActivityId,
|
||||
notification.WorkflowExecutionContext.WorkflowInstance.TenantId,
|
||||
cancellationToken);
|
||||
}
|
||||
|
||||
public Task Handle(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) => _workflowScheduler.UnscheduleWorkflowDefinitionAsync(notification.WorkflowDefinition.DefinitionId, notification.WorkflowDefinition.TenantId, cancellationToken);
|
||||
public Task Handle(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => _workflowScheduler.UnscheduleWorkflowDefinitionAsync(notification.WorkflowDefinition.DefinitionId, notification.WorkflowDefinition.TenantId, cancellationToken);
|
||||
public Task Handle(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) => _workflowDefinitionScheduler.UnscheduleAsync(notification.WorkflowDefinition.DefinitionId , cancellationToken);
|
||||
public Task Handle(WorkflowDefinitionRetracted notification, CancellationToken cancellationToken) => _workflowDefinitionScheduler.UnscheduleAsync(notification.WorkflowDefinition.DefinitionId, cancellationToken);
|
||||
public Task Handle(WorkflowDefinitionDeleted notification, CancellationToken cancellationToken) => _workflowDefinitionScheduler.UnscheduleAsync(notification.WorkflowDefinition.DefinitionId, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -2,6 +2,7 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.Temporal.Common.Bookmarks;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.Activities.Temporal.Common.Services;
|
||||
using Elsa.Events;
|
||||
using Elsa.Triggers;
|
||||
|
|
@ -15,11 +16,11 @@ namespace Elsa.Activities.Temporal.Common.Handlers
|
|||
// TODO: Design multi-tenancy.
|
||||
private const string? TenantId = default;
|
||||
|
||||
private readonly IWorkflowScheduler _workflowScheduler;
|
||||
private readonly IWorkflowDefinitionScheduler _workflowScheduler;
|
||||
private readonly ITriggerFinder _triggerFinder;
|
||||
private readonly ILogger<SchedulePublishedWorkflows> _logger;
|
||||
|
||||
public SchedulePublishedWorkflows(IWorkflowScheduler workflowScheduler, ITriggerFinder triggerFinder, ILogger<SchedulePublishedWorkflows> logger)
|
||||
public SchedulePublishedWorkflows(IWorkflowDefinitionScheduler workflowScheduler, ITriggerFinder triggerFinder, ILogger<SchedulePublishedWorkflows> logger)
|
||||
{
|
||||
_workflowScheduler = workflowScheduler;
|
||||
_triggerFinder = triggerFinder;
|
||||
|
|
@ -35,19 +36,19 @@ namespace Elsa.Activities.Temporal.Common.Handlers
|
|||
foreach (var trigger in startAtTriggers)
|
||||
{
|
||||
var bookmark = (StartAtBookmark) trigger.Bookmark;
|
||||
await Try(() => _workflowScheduler.ScheduleWorkflowAsync(trigger.WorkflowBlueprint.Id, null, trigger.ActivityId, TenantId, bookmark.ExecuteAt, null, cancellationToken));
|
||||
await Try(() => _workflowScheduler.ScheduleAsync(trigger.WorkflowBlueprint.Id, trigger.ActivityId, bookmark.ExecuteAt, null, ClusterMode.SingleNode, cancellationToken));
|
||||
}
|
||||
|
||||
foreach (var trigger in timerTriggers)
|
||||
{
|
||||
var bookmark = (TimerBookmark) trigger.Bookmark;
|
||||
await Try(() => _workflowScheduler.ScheduleWorkflowAsync(trigger.WorkflowBlueprint.Id, null, trigger.ActivityId, TenantId, bookmark.ExecuteAt, bookmark.Interval, cancellationToken));
|
||||
await Try(() => _workflowScheduler.ScheduleAsync(trigger.WorkflowBlueprint.Id, trigger.ActivityId, bookmark.ExecuteAt, bookmark.Interval, bookmark.ClusterMode ?? ClusterMode.SingleNode, cancellationToken));
|
||||
}
|
||||
|
||||
foreach (var trigger in cronTriggers)
|
||||
{
|
||||
var bookmark = (CronBookmark) trigger.Bookmark;
|
||||
await Try(() => _workflowScheduler.ScheduleWorkflowAsync(trigger.WorkflowBlueprint.Id, null, trigger.ActivityId, TenantId, bookmark.CronExpression, cancellationToken));
|
||||
await Try(() => _workflowScheduler.ScheduleAsync(trigger.WorkflowBlueprint.Id, trigger.ActivityId, bookmark.CronExpression, ClusterMode.SingleNode, cancellationToken));
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,8 @@
|
|||
namespace Elsa.Activities.Temporal.Common.Options
|
||||
{
|
||||
public enum ClusterMode
|
||||
{
|
||||
SingleNode,
|
||||
MultiNode
|
||||
}
|
||||
}
|
||||
|
|
@ -1,14 +1,25 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using NodaTime;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Common.Services
|
||||
{
|
||||
public interface IWorkflowScheduler
|
||||
public interface IWorkflowDefinitionScheduler
|
||||
{
|
||||
Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, Instant startAt, Duration? interval = default, CancellationToken cancellationToken = default);
|
||||
Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string cronExpression, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleWorkflowDefinitionAsync(string workflowDefinitionId, string? tenantId, CancellationToken cancellationToken = default);
|
||||
Task ScheduleAsync(string workflowDefinitionId, string activityId, Instant startAt, Duration? interval, ClusterMode clusterMode = ClusterMode.SingleNode, CancellationToken cancellationToken = default);
|
||||
Task ScheduleAsync(string workflowDefinitionId, string activityId, string cronExpression, ClusterMode clusterMode = ClusterMode.SingleNode, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleAsync(string workflowDefinitionId, string activityId, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleAsync(string workflowDefinitionId, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleAllAsync(CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
||||
public interface IWorkflowInstanceScheduler
|
||||
{
|
||||
Task ScheduleAsync(string workflowInstanceId, string activityId, Instant startAt, Duration? interval, CancellationToken cancellationToken = default);
|
||||
Task ScheduleAsync(string workflowInstanceId, string activityId, string cronExpression, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleAsync(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleAsync(string workflowInstanceId, CancellationToken cancellationToken = default);
|
||||
Task UnscheduleAllAsync(CancellationToken cancellationToken = default);
|
||||
}
|
||||
}
|
||||
|
|
@ -16,9 +16,9 @@ namespace Elsa.Activities.Temporal.Common.StartupTasks
|
|||
private const string? TenantId = default;
|
||||
|
||||
private readonly IBookmarkFinder _bookmarkFinder;
|
||||
private readonly IWorkflowScheduler _workflowScheduler;
|
||||
private readonly IWorkflowInstanceScheduler _workflowScheduler;
|
||||
|
||||
public StartJobs(IBookmarkFinder bookmarkFinder, IWorkflowScheduler workflowScheduler)
|
||||
public StartJobs(IBookmarkFinder bookmarkFinder, IWorkflowInstanceScheduler workflowScheduler)
|
||||
{
|
||||
_bookmarkFinder = bookmarkFinder;
|
||||
_workflowScheduler = workflowScheduler;
|
||||
|
|
@ -41,7 +41,7 @@ namespace Elsa.Activities.Temporal.Common.StartupTasks
|
|||
foreach (var result in bookmarkResults)
|
||||
{
|
||||
var bookmark = (StartAtBookmark) result.Bookmark;
|
||||
await _workflowScheduler.ScheduleWorkflowAsync(null, result.WorkflowInstanceId!, result.ActivityId, TenantId, bookmark.ExecuteAt, null, cancellationToken);
|
||||
await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, bookmark.ExecuteAt, null, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -53,7 +53,7 @@ namespace Elsa.Activities.Temporal.Common.StartupTasks
|
|||
foreach (var result in bookmarkResults)
|
||||
{
|
||||
var bookmark = (TimerBookmark) result.Bookmark;
|
||||
await _workflowScheduler.ScheduleWorkflowAsync(null, result.WorkflowInstanceId!, result.ActivityId, TenantId, bookmark.ExecuteAt, null, cancellationToken);
|
||||
await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, bookmark.ExecuteAt, null, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -65,7 +65,7 @@ namespace Elsa.Activities.Temporal.Common.StartupTasks
|
|||
foreach (var result in cronEventTriggers)
|
||||
{
|
||||
var trigger = (CronBookmark) result.Bookmark;
|
||||
await _workflowScheduler.ScheduleWorkflowAsync(null, result.WorkflowInstanceId!, result.ActivityId, TenantId, trigger.ExecuteAt!.Value, null, cancellationToken);
|
||||
await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, trigger.ExecuteAt!.Value, null, cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,18 +1,20 @@
|
|||
using System;
|
||||
using System.Security.Cryptography;
|
||||
using System.Text;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Hangfire.Models
|
||||
{
|
||||
public class RunHangfireWorkflowJobModel
|
||||
{
|
||||
public RunHangfireWorkflowJobModel(string? workflowDefinitionId, string activityId, string? workflowInstanceId, string? tenantId, string? cronExpression)
|
||||
public RunHangfireWorkflowJobModel(string? workflowDefinitionId, string activityId, string? workflowInstanceId, string? tenantId, string? cronExpression, ClusterMode? clusterMode)
|
||||
{
|
||||
WorkflowDefinitionId = workflowDefinitionId;
|
||||
WorkflowInstanceId = workflowInstanceId;
|
||||
ActivityId = activityId;
|
||||
TenantId = tenantId;
|
||||
CronExpression = cronExpression;
|
||||
ClusterMode = clusterMode;
|
||||
}
|
||||
|
||||
public string? WorkflowDefinitionId { get; set; }
|
||||
|
|
@ -20,6 +22,7 @@ namespace Elsa.Activities.Temporal.Hangfire.Models
|
|||
public string ActivityId { get; set; }
|
||||
public string? TenantId { get; set; }
|
||||
public string? CronExpression { get; set; }
|
||||
public ClusterMode? ClusterMode { get; }
|
||||
public bool IsRecurringJob => string.IsNullOrEmpty(CronExpression) == false;
|
||||
|
||||
public string GetIdentity()
|
||||
|
|
|
|||
|
|
@ -1,9 +1,9 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.Activities.Temporal.Common.Services;
|
||||
using Elsa.Activities.Temporal.Hangfire.Extensions;
|
||||
using Elsa.Activities.Temporal.Hangfire.Models;
|
||||
using Hangfire;
|
||||
using NodaTime;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Hangfire.Services
|
||||
|
|
@ -17,22 +17,22 @@ namespace Elsa.Activities.Temporal.Hangfire.Services
|
|||
_jobManager = jobManager;
|
||||
}
|
||||
|
||||
public Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, Instant startAt, Duration? interval, CancellationToken cancellationToken)
|
||||
public Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, Instant startAt, Duration? interval, ClusterMode? clusterMode = default, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var cron = interval?.ToCronExpression();
|
||||
var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, cron);
|
||||
var cronExpression = interval?.ToCronExpression();
|
||||
var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, cronExpression, clusterMode);
|
||||
|
||||
_jobManager.ScheduleJob(data, startAt);
|
||||
|
||||
if (cron != null)
|
||||
_jobManager.ScheduleRecurringJob(data, cron);
|
||||
if (cronExpression != null)
|
||||
_jobManager.ScheduleRecurringJob(data, cronExpression);
|
||||
|
||||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
public Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string cronExpression, CancellationToken cancellationToken)
|
||||
public Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string cronExpression, ClusterMode? clusterMode = default, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, cronExpression);
|
||||
var data = CreateData(workflowDefinitionId, workflowInstanceId, activityId, tenantId, cronExpression, clusterMode);
|
||||
|
||||
_jobManager.ScheduleRecurringJob(data, cronExpression);
|
||||
|
||||
|
|
@ -52,14 +52,15 @@ namespace Elsa.Activities.Temporal.Hangfire.Services
|
|||
return Task.CompletedTask;
|
||||
}
|
||||
|
||||
private RunHangfireWorkflowJobModel CreateData(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string? cronExpression = null)
|
||||
private RunHangfireWorkflowJobModel CreateData(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string? cronExpression = default, ClusterMode? clusterMode = default)
|
||||
{
|
||||
return new(
|
||||
workflowDefinitionId,
|
||||
workflowInstanceId: workflowInstanceId,
|
||||
activityId: activityId,
|
||||
tenantId: tenantId,
|
||||
cronExpression: cronExpression);
|
||||
cronExpression: cronExpression,
|
||||
clusterMode: clusterMode);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -28,9 +28,10 @@ namespace Elsa
|
|||
Action<QuartzHostedServiceOptions>? configureQuartzHostedService = default)
|
||||
{
|
||||
timersOptions.Services
|
||||
.AddSingleton<IWorkflowScheduler, QuartzWorkflowScheduler>()
|
||||
.AddSingleton<ICrontabParser, QuartzCrontabParser>()
|
||||
.AddTransient<RunQuartzWorkflowJob>()
|
||||
.AddSingleton<IWorkflowDefinitionScheduler, QuartzWorkflowDefinitionScheduler>()
|
||||
.AddSingleton<IWorkflowInstanceScheduler, QuartzWorkflowInstanceScheduler>()
|
||||
.AddTransient<RunQuartzWorkflowDefinitionJob>()
|
||||
.AddTransient<RunQuartzWorkflowInstanceJob>()
|
||||
.AddNotificationHandlers(typeof(ConfigureCronProperty));
|
||||
|
||||
if (registerQuartz)
|
||||
|
|
@ -52,11 +53,11 @@ namespace Elsa
|
|||
|
||||
private static void ConfigureQuartz(IServiceCollectionQuartzConfigurator quartz, Action<IServiceCollectionQuartzConfigurator>? configureQuartz)
|
||||
{
|
||||
quartz.UseMicrosoftDependencyInjectionScopedJobFactory(options => options.AllowDefaultConstructor = true);
|
||||
quartz.AddJob<RunQuartzWorkflowJob>(job => job.StoreDurably().WithIdentity(nameof(RunQuartzWorkflowJob)));
|
||||
quartz.UseMicrosoftDependencyInjectionScopedJobFactory();
|
||||
quartz.AddJob<RunQuartzWorkflowDefinitionJob>(job => job.StoreDurably().WithIdentity(nameof(RunQuartzWorkflowDefinitionJob)));
|
||||
quartz.AddJob<RunQuartzWorkflowInstanceJob>(job => job.StoreDurably().WithIdentity(nameof(RunQuartzWorkflowInstanceJob)));
|
||||
quartz.UseSimpleTypeLoader();
|
||||
quartz.UseInMemoryStore();
|
||||
|
||||
configureQuartz?.Invoke(quartz);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,48 @@
|
|||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.Dispatch;
|
||||
using Elsa.Services;
|
||||
using Quartz;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Quartz.Jobs
|
||||
{
|
||||
public class RunQuartzWorkflowDefinitionJob : IJob
|
||||
{
|
||||
private readonly IWorkflowDefinitionDispatcher _workflowDefinitionDispatcher;
|
||||
private readonly IDistributedLockProvider _distributedLockProvider;
|
||||
|
||||
public RunQuartzWorkflowDefinitionJob(
|
||||
IWorkflowDefinitionDispatcher workflowDefinitionDispatcher,
|
||||
IDistributedLockProvider distributedLockProvider)
|
||||
{
|
||||
_workflowDefinitionDispatcher = workflowDefinitionDispatcher;
|
||||
_distributedLockProvider = distributedLockProvider;
|
||||
}
|
||||
|
||||
public async Task Execute(IJobExecutionContext context)
|
||||
{
|
||||
var dataMap = context.MergedJobDataMap;
|
||||
var cancellationToken = context.CancellationToken;
|
||||
var workflowDefinitionId = dataMap.GetString("WorkflowDefinitionId")!;
|
||||
var activityId = dataMap.GetString("ActivityId")!;
|
||||
var clusterModeString = dataMap.GetString("ClusterMode");
|
||||
var clusterMode = !string.IsNullOrEmpty(clusterModeString) ? Enum.Parse<ClusterMode>(clusterModeString) : ClusterMode.SingleNode;
|
||||
|
||||
if (clusterMode == ClusterMode.MultiNode)
|
||||
{
|
||||
await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowDefinitionId, activityId), cancellationToken);
|
||||
}
|
||||
else
|
||||
{
|
||||
var resource = $"{nameof(RunQuartzWorkflowDefinitionJob)}:{workflowDefinitionId}";
|
||||
await using var @lock = await _distributedLockProvider.AcquireLockAsync(resource, cancellationToken: cancellationToken);
|
||||
|
||||
if (@lock == null)
|
||||
return;
|
||||
|
||||
await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowDefinitionId, activityId), cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,26 @@
|
|||
using System.Threading.Tasks;
|
||||
using Elsa.Dispatch;
|
||||
using Quartz;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Quartz.Jobs
|
||||
{
|
||||
public class RunQuartzWorkflowInstanceJob : IJob
|
||||
{
|
||||
private readonly IWorkflowInstanceDispatcher _workflowInstanceDispatcher;
|
||||
|
||||
public RunQuartzWorkflowInstanceJob(IWorkflowInstanceDispatcher workflowDefinitionDispatcher)
|
||||
{
|
||||
_workflowInstanceDispatcher = workflowDefinitionDispatcher;
|
||||
}
|
||||
|
||||
public async Task Execute(IJobExecutionContext context)
|
||||
{
|
||||
var dataMap = context.MergedJobDataMap;
|
||||
var cancellationToken = context.CancellationToken;
|
||||
var workflowInstanceId = dataMap.GetString("WorkflowInstanceId")!;
|
||||
var activityId = dataMap.GetString("ActivityId")!;
|
||||
|
||||
await _workflowInstanceDispatcher.DispatchAsync(new ExecuteWorkflowInstanceRequest(workflowInstanceId, activityId), cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,35 +0,0 @@
|
|||
using System.Threading.Tasks;
|
||||
using Elsa.Dispatch;
|
||||
using Quartz;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Quartz.Jobs
|
||||
{
|
||||
public class RunQuartzWorkflowJob : IJob
|
||||
{
|
||||
private readonly IWorkflowDefinitionDispatcher _workflowDefinitionDispatcher;
|
||||
private readonly IWorkflowInstanceDispatcher _workflowInstanceDispatcher;
|
||||
|
||||
public RunQuartzWorkflowJob(
|
||||
IWorkflowDefinitionDispatcher workflowDefinitionDispatcher,
|
||||
IWorkflowInstanceDispatcher workflowInstanceDispatcher)
|
||||
{
|
||||
_workflowDefinitionDispatcher = workflowDefinitionDispatcher;
|
||||
_workflowInstanceDispatcher = workflowInstanceDispatcher;
|
||||
}
|
||||
|
||||
public async Task Execute(IJobExecutionContext context)
|
||||
{
|
||||
var dataMap = context.MergedJobDataMap;
|
||||
var cancellationToken = context.CancellationToken;
|
||||
var workflowInstanceId = dataMap.GetString("WorkflowInstanceId");
|
||||
var tenantId = dataMap.GetString("TenantId");
|
||||
var workflowDefinitionId = dataMap.GetString("WorkflowDefinitionId")!;
|
||||
var activityId = dataMap.GetString("ActivityId")!;
|
||||
|
||||
if (workflowInstanceId == null)
|
||||
await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowDefinitionId, activityId, TenantId: tenantId), cancellationToken);
|
||||
else
|
||||
await _workflowInstanceDispatcher.DispatchAsync(new ExecuteWorkflowInstanceRequest(workflowInstanceId, activityId), cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,110 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.Activities.Temporal.Common.Services;
|
||||
using Elsa.Activities.Temporal.Quartz.Jobs;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using NodaTime;
|
||||
using Quartz;
|
||||
using Quartz.Impl.Matchers;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Quartz.Services
|
||||
{
|
||||
public class QuartzWorkflowDefinitionScheduler : IWorkflowDefinitionScheduler
|
||||
{
|
||||
private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowDefinitionJob);
|
||||
private readonly ISchedulerFactory _schedulerFactory;
|
||||
private readonly ILogger _logger;
|
||||
private readonly SemaphoreSlim _semaphore = new(1);
|
||||
|
||||
public QuartzWorkflowDefinitionScheduler(ISchedulerFactory schedulerFactory, ILogger<QuartzWorkflowDefinitionScheduler> logger)
|
||||
{
|
||||
_schedulerFactory = schedulerFactory;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
public async Task ScheduleAsync(string workflowDefinitionId, string activityId, Instant startAt, Duration? interval, ClusterMode clusterMode, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var triggerBuilder = CreateTrigger(workflowDefinitionId, activityId, clusterMode).StartAt(startAt.ToDateTimeOffset());
|
||||
|
||||
if (interval != null && interval != Duration.Zero)
|
||||
triggerBuilder.WithSimpleSchedule(x => x.WithInterval(interval.Value.ToTimeSpan()).RepeatForever());
|
||||
|
||||
var trigger = triggerBuilder.Build();
|
||||
await ScheduleJob(trigger, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task ScheduleAsync(string workflowDefinitionId, string activityId, string cronExpression, ClusterMode clusterMode, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var trigger = CreateTrigger(workflowDefinitionId, activityId, clusterMode).WithCronSchedule(cronExpression).Build();
|
||||
await ScheduleJob(trigger, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UnscheduleAsync(string workflowDefinitionId, string activityId, CancellationToken cancellationToken)
|
||||
{
|
||||
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
|
||||
var trigger = CreateTriggerKey(workflowDefinitionId, activityId);
|
||||
var existingTrigger = await scheduler.GetTrigger(trigger, cancellationToken);
|
||||
|
||||
if (existingTrigger != null)
|
||||
await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UnscheduleAsync(string workflowDefinitionId, CancellationToken cancellationToken)
|
||||
{
|
||||
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
|
||||
var groupName = CreateTriggerGroupKey(workflowDefinitionId);
|
||||
var existingTriggers = await scheduler.GetTriggerKeys(GroupMatcher<TriggerKey>.GroupEquals(groupName), cancellationToken);
|
||||
|
||||
foreach (var existingTrigger in existingTriggers)
|
||||
await scheduler.UnscheduleJob(existingTrigger, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UnscheduleAllAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
|
||||
var jobKeys = await scheduler.GetJobKeys(GroupMatcher<JobKey>.GroupStartsWith("workflow-definition"), cancellationToken);
|
||||
await scheduler.DeleteJobs(jobKeys, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken)
|
||||
{
|
||||
await _semaphore.WaitAsync(cancellationToken);
|
||||
|
||||
try
|
||||
{
|
||||
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
|
||||
var existingTrigger = await scheduler.GetTrigger(trigger.Key, cancellationToken);
|
||||
|
||||
if (existingTrigger != null)
|
||||
await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken);
|
||||
|
||||
await scheduler.ScheduleJob(trigger, cancellationToken);
|
||||
}
|
||||
catch (SchedulerException e)
|
||||
{
|
||||
_logger.LogWarning(e, "Failed to schedule trigger {TriggerKey}", trigger.Key.ToString());
|
||||
}
|
||||
finally
|
||||
{
|
||||
_semaphore.Release();
|
||||
}
|
||||
}
|
||||
|
||||
private TriggerBuilder CreateTrigger(string workflowDefinitionId, string activityId, ClusterMode clusterMode) =>
|
||||
TriggerBuilder.Create()
|
||||
.ForJob(RunWorkflowJobKey)
|
||||
.WithIdentity(CreateTriggerKey(workflowDefinitionId, activityId))
|
||||
.UsingJobData("WorkflowDefinitionId", workflowDefinitionId!)
|
||||
.UsingJobData("ActivityId", activityId)
|
||||
.UsingJobData("ClusterMode", clusterMode.ToString());
|
||||
|
||||
private TriggerKey CreateTriggerKey(string workflowDefinitionId, string activityId)
|
||||
{
|
||||
var groupName = CreateTriggerGroupKey(workflowDefinitionId);
|
||||
return new TriggerKey($"activity:{activityId}", groupName);
|
||||
}
|
||||
|
||||
private string CreateTriggerGroupKey(string workflowDefinitionId) => $"workflow-definition:{workflowDefinitionId}";
|
||||
}
|
||||
}
|
||||
|
|
@ -1,5 +1,6 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Activities.Temporal.Common.Options;
|
||||
using Elsa.Activities.Temporal.Common.Services;
|
||||
using Elsa.Activities.Temporal.Quartz.Jobs;
|
||||
using Microsoft.Extensions.Logging;
|
||||
|
|
@ -9,22 +10,22 @@ using Quartz.Impl.Matchers;
|
|||
|
||||
namespace Elsa.Activities.Temporal.Quartz.Services
|
||||
{
|
||||
public class QuartzWorkflowScheduler : IWorkflowScheduler
|
||||
public class QuartzWorkflowInstanceScheduler : IWorkflowInstanceScheduler
|
||||
{
|
||||
private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowJob);
|
||||
private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowInstanceJob);
|
||||
private readonly ISchedulerFactory _schedulerFactory;
|
||||
private readonly ILogger _logger;
|
||||
private readonly SemaphoreSlim _semaphore = new(1);
|
||||
|
||||
public QuartzWorkflowScheduler(ISchedulerFactory schedulerFactory, ILogger<QuartzWorkflowScheduler> logger)
|
||||
public QuartzWorkflowInstanceScheduler(ISchedulerFactory schedulerFactory, ILogger<QuartzWorkflowInstanceScheduler> logger)
|
||||
{
|
||||
_schedulerFactory = schedulerFactory;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
public async Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, Instant startAt, Duration? interval, CancellationToken cancellationToken)
|
||||
public async Task ScheduleAsync(string workflowInstanceId, string activityId, Instant startAt, Duration? interval, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var triggerBuilder = CreateTrigger(workflowDefinitionId, workflowInstanceId, activityId, tenantId).StartAt(startAt.ToDateTimeOffset());
|
||||
var triggerBuilder = CreateTrigger(workflowInstanceId, activityId).StartAt(startAt.ToDateTimeOffset());
|
||||
|
||||
if (interval != null && interval != Duration.Zero)
|
||||
triggerBuilder.WithSimpleSchedule(x => x.WithInterval(interval.Value.ToTimeSpan()).RepeatForever());
|
||||
|
|
@ -33,32 +34,39 @@ namespace Elsa.Activities.Temporal.Quartz.Services
|
|||
await ScheduleJob(trigger, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task ScheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, string cronExpression, CancellationToken cancellationToken)
|
||||
public async Task ScheduleAsync(string workflowInstanceId, string activityId, string cronExpression, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var trigger = CreateTrigger(workflowDefinitionId, workflowInstanceId, activityId, tenantId).WithCronSchedule(cronExpression).Build();
|
||||
var trigger = CreateTrigger(workflowInstanceId, activityId).WithCronSchedule(cronExpression).Build();
|
||||
await ScheduleJob(trigger, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UnscheduleWorkflowAsync(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId, CancellationToken cancellationToken)
|
||||
public async Task UnscheduleAsync(string workflowInstanceId, string activityId, CancellationToken cancellationToken)
|
||||
{
|
||||
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
|
||||
var trigger = CreateTriggerKey(tenantId, workflowDefinitionId, workflowInstanceId, activityId);
|
||||
var trigger = CreateTriggerKey(workflowInstanceId, activityId);
|
||||
var existingTrigger = await scheduler.GetTrigger(trigger, cancellationToken);
|
||||
|
||||
if (existingTrigger != null)
|
||||
await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UnscheduleWorkflowDefinitionAsync(string workflowDefinitionId, string? tenantId, CancellationToken cancellationToken)
|
||||
public async Task UnscheduleAsync(string workflowInstanceId, CancellationToken cancellationToken)
|
||||
{
|
||||
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
|
||||
var groupName = CreateTriggerGroupKey(tenantId, workflowDefinitionId);
|
||||
var groupName = CreateTriggerGroupKey(workflowInstanceId);
|
||||
var existingTriggers = await scheduler.GetTriggerKeys(GroupMatcher<TriggerKey>.GroupEquals(groupName), cancellationToken);
|
||||
|
||||
foreach (var existingTrigger in existingTriggers)
|
||||
await scheduler.UnscheduleJob(existingTrigger, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task UnscheduleAllAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
var scheduler = await _schedulerFactory.GetScheduler(cancellationToken);
|
||||
var jobKeys = await scheduler.GetJobKeys(GroupMatcher<JobKey>.GroupStartsWith("workflow-instance"), cancellationToken);
|
||||
await scheduler.DeleteJobs(jobKeys, cancellationToken);
|
||||
}
|
||||
|
||||
private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken)
|
||||
{
|
||||
await _semaphore.WaitAsync(cancellationToken);
|
||||
|
|
@ -83,21 +91,19 @@ namespace Elsa.Activities.Temporal.Quartz.Services
|
|||
}
|
||||
}
|
||||
|
||||
private TriggerBuilder CreateTrigger(string? workflowDefinitionId, string? workflowInstanceId, string activityId, string? tenantId) =>
|
||||
private TriggerBuilder CreateTrigger(string workflowInstanceId, string activityId) =>
|
||||
TriggerBuilder.Create()
|
||||
.ForJob(RunWorkflowJobKey)
|
||||
.WithIdentity(CreateTriggerKey(tenantId, workflowDefinitionId, workflowInstanceId, activityId))
|
||||
.UsingJobData("TenantId", tenantId!)
|
||||
.UsingJobData("WorkflowDefinitionId", workflowDefinitionId!)
|
||||
.WithIdentity(CreateTriggerKey(workflowInstanceId, activityId))
|
||||
.UsingJobData("WorkflowInstanceId", workflowInstanceId!)
|
||||
.UsingJobData("ActivityId", activityId);
|
||||
|
||||
private TriggerKey CreateTriggerKey(string? tenantId, string? workflowDefinitionId, string? workflowInstanceId, string activityId)
|
||||
|
||||
private TriggerKey CreateTriggerKey(string workflowInstanceId, string activityId)
|
||||
{
|
||||
var groupName = CreateTriggerGroupKey(tenantId, workflowInstanceId ?? workflowDefinitionId);
|
||||
var groupName = CreateTriggerGroupKey(workflowInstanceId);
|
||||
return new TriggerKey($"activity:{activityId}", groupName);
|
||||
}
|
||||
|
||||
private string CreateTriggerGroupKey(string? tenantId, string? workflowId) => $"tenant:{tenantId ?? "default"}-workflow:{workflowId}";
|
||||
private string CreateTriggerGroupKey(string workflowInstanceId) => $"workflow-instance:{workflowInstanceId}";
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,9 @@
|
|||
using System;
|
||||
|
||||
namespace Elsa.Attributes
|
||||
{
|
||||
[AttributeUsage(AttributeTargets.Property)]
|
||||
public class ExcludeFromHashAttribute : Attribute
|
||||
{
|
||||
}
|
||||
}
|
||||
|
|
@ -3,5 +3,6 @@
|
|||
public static class PropertyCategories
|
||||
{
|
||||
public const string Advanced = "Advanced";
|
||||
public const string Hosting = "Hosting";
|
||||
}
|
||||
}
|
||||
|
|
@ -1,5 +1,10 @@
|
|||
using System.Security.Cryptography;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Reflection;
|
||||
using System.Security.Cryptography;
|
||||
using System.Text;
|
||||
using Elsa.Attributes;
|
||||
using Elsa.Serialization.ContractResolvers;
|
||||
using Newtonsoft.Json;
|
||||
using NodaTime;
|
||||
using NodaTime.Serialization.JsonNet;
|
||||
|
|
@ -8,17 +13,19 @@ namespace Elsa.Bookmarks
|
|||
{
|
||||
public class BookmarkHasher : IBookmarkHasher
|
||||
{
|
||||
private readonly JsonSerializerSettings _serializerSettings;
|
||||
|
||||
public BookmarkHasher()
|
||||
{
|
||||
_serializerSettings = new JsonSerializerSettings().ConfigureForNodaTime(DateTimeZoneProviders.Tzdb);
|
||||
}
|
||||
|
||||
public string Hash(IBookmark bookmark)
|
||||
{
|
||||
var json = JsonConvert.SerializeObject(bookmark, _serializerSettings);
|
||||
var hash = Hash(json);
|
||||
var type = bookmark.GetType();
|
||||
var whiteListedProperties = FilterProperties(type.GetProperties()).ToList();
|
||||
var contractResolver = new WhiteListedPropertiesContractResolver(whiteListedProperties);
|
||||
|
||||
var serializerSettings = new JsonSerializerSettings
|
||||
{
|
||||
ContractResolver = contractResolver
|
||||
}.ConfigureForNodaTime(DateTimeZoneProviders.Tzdb);
|
||||
|
||||
var json = JsonConvert.SerializeObject(bookmark, serializerSettings);
|
||||
var hash = Hash(json);
|
||||
return hash;
|
||||
}
|
||||
|
||||
|
|
@ -38,5 +45,7 @@ namespace Elsa.Bookmarks
|
|||
|
||||
return builder.ToString();
|
||||
}
|
||||
|
||||
private static IEnumerable<PropertyInfo> FilterProperties(IEnumerable<PropertyInfo> properties) => properties.Where(property => property.GetCustomAttribute<ExcludeFromHashAttribute>() == null);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,26 @@
|
|||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Reflection;
|
||||
using Newtonsoft.Json;
|
||||
using Newtonsoft.Json.Serialization;
|
||||
|
||||
namespace Elsa.Serialization.ContractResolvers
|
||||
{
|
||||
public class WhiteListedPropertiesContractResolver : DefaultContractResolver
|
||||
{
|
||||
private readonly string[] _whiteList;
|
||||
public WhiteListedPropertiesContractResolver(params string[] whiteList) => _whiteList = whiteList;
|
||||
|
||||
public WhiteListedPropertiesContractResolver(IEnumerable<PropertyInfo> properties) : this(properties.Select(x => x.Name).ToArray())
|
||||
{
|
||||
}
|
||||
|
||||
protected override IList<JsonProperty> CreateProperties(Type type, MemberSerialization memberSerialization)
|
||||
{
|
||||
var props = base.CreateProperties(type, memberSerialization);
|
||||
props = props.Where(p => _whiteList.Contains(p.PropertyName)).ToList();
|
||||
return props;
|
||||
}
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large
Load diff
File diff suppressed because it is too large
Load diff
|
|
@ -34,6 +34,7 @@
|
|||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Hangfire.InMemory" Version="0.3.4" />
|
||||
<PackageReference Include="Hangfire.SqlServer" Version="1.7.22" />
|
||||
<PackageReference Include="Quartz.Serialization.Json" Version="3.3.2" />
|
||||
</ItemGroup>
|
||||
|
||||
|
|
|
|||
|
|
@ -1,12 +1,17 @@
|
|||
using System;
|
||||
using Elsa.Activities.UserTask.Extensions;
|
||||
using Elsa.Caching.Rebus.Extensions;
|
||||
using Elsa.Persistence.EntityFramework.Core.Extensions;
|
||||
using Elsa.Persistence.EntityFramework.Sqlite;
|
||||
using Elsa.Persistence.EntityFramework.SqlServer;
|
||||
using Elsa.Rebus.RabbitMq;
|
||||
using Elsa.Samples.Server.Host.Activities;
|
||||
using Hangfire;
|
||||
using Microsoft.AspNetCore.Builder;
|
||||
using Microsoft.AspNetCore.Hosting;
|
||||
using Microsoft.Extensions.Configuration;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Quartz;
|
||||
|
||||
namespace Elsa.Samples.Server.Host
|
||||
{
|
||||
|
|
@ -24,6 +29,7 @@ namespace Elsa.Samples.Server.Host
|
|||
public void ConfigureServices(IServiceCollection services)
|
||||
{
|
||||
var elsaSection = Configuration.GetSection("Elsa");
|
||||
var sqlServerConnectionString = Configuration.GetConnectionString("SqlServer");
|
||||
|
||||
services.AddControllers();
|
||||
|
||||
|
|
@ -32,21 +38,23 @@ namespace Elsa.Samples.Server.Host
|
|||
.AddRuntimeSelectItemsProvider<VehicleActivity>()
|
||||
//.AddRedis(Configuration.GetConnectionString("Redis"))
|
||||
.AddElsa(elsa => elsa
|
||||
//.WithContainerName(Configuration["ContainerName"] ?? System.Environment.MachineName)
|
||||
.UseEntityFrameworkPersistence(ef => ef.UseSqlite())
|
||||
//.UseRabbitMq(Configuration.GetConnectionString("RabbitMq"))
|
||||
//.UseRebusCacheSignal()
|
||||
.WithContainerName(Configuration["ContainerName"] ?? System.Environment.MachineName)
|
||||
.UseEntityFrameworkPersistence(ef => ef.UseSqlServer(sqlServerConnectionString))
|
||||
.UseRabbitMq(Configuration.GetConnectionString("RabbitMq"))
|
||||
.UseRebusCacheSignal()
|
||||
//.UseRedisCacheSignal()
|
||||
// .ConfigureDistributedLockProvider(options => options.UseProviderFactory(sp => name =>
|
||||
// {
|
||||
// var connection = sp.GetRequiredService<IConnectionMultiplexer>();
|
||||
// return new RedisDistributedLock(name, connection.GetDatabase());
|
||||
// }))
|
||||
.AddConsoleActivities()
|
||||
.AddHttpActivities(elsaSection.GetSection("Http").Bind)
|
||||
.AddEmailActivities(elsaSection.GetSection("Smtp").Bind)
|
||||
.AddQuartzTemporalActivities()
|
||||
// .AddQuartzTemporalActivities(configureQuartz: quartz => quartz.UsePersistentStore(store =>
|
||||
// {
|
||||
// store.UseJsonSerializer();
|
||||
// store.UseSqlServer(sqlServerConnectionString);
|
||||
// store.UseClustering();
|
||||
// }))
|
||||
//.AddHangfireTemporalActivities(hangfire => hangfire.UseInMemoryStorage(), (_, hangfireServer) => hangfireServer.SchedulePollingInterval = TimeSpan.FromSeconds(5))
|
||||
//.AddHangfireTemporalActivities(hangfire => hangfire.UseSqlServerStorage(sqlServerConnectionString), (_, hangfireServer) => hangfireServer.SchedulePollingInterval = TimeSpan.FromSeconds(5))
|
||||
.AddJavaScriptActivities()
|
||||
.AddUserTaskActivities()
|
||||
.AddActivitiesFrom<Startup>()
|
||||
|
|
|
|||
|
|
@ -1,7 +1,7 @@
|
|||
{
|
||||
"Logging": {
|
||||
"LogLevel": {
|
||||
"Default": "Debug",
|
||||
"Default": "Warning",
|
||||
"Orleans": "Warning",
|
||||
"System": "Information",
|
||||
"Microsoft": "Warning"
|
||||
|
|
|
|||
|
|
@ -10,7 +10,8 @@
|
|||
"AllowedHosts": "*",
|
||||
"ConnectionStrings": {
|
||||
"RabbitMq": "amqp://localhost:5672",
|
||||
"Redis": "localhost:6379,abortConnect=false"
|
||||
"Redis": "localhost:6379,abortConnect=false",
|
||||
"SqlServer": "Server=LAPTOP-B76STK67;Database=Elsa;Integrated Security=true;MultipleActiveResultSets=True;Max Pool Size=500;Connection Timeout=3600"
|
||||
},
|
||||
"Elsa": {
|
||||
"Http": {
|
||||
|
|
|
|||
26
src/server/Elsa.Server.Host/Elsa.Server.Host.csproj
Normal file
26
src/server/Elsa.Server.Host/Elsa.Server.Host.csproj
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
<Project Sdk="Microsoft.NET.Sdk.Web">
|
||||
<Import Project="..\..\..\configureawait.props" />
|
||||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>net5.0</TargetFramework>
|
||||
<LangVersion>latest</LangVersion>
|
||||
<IsPackable>false</IsPackable>
|
||||
<Nullable>enable</Nullable>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Hangfire.InMemory" Version="0.3.4" />
|
||||
<PackageReference Include="Hangfire.SqlServer" Version="1.7.22" />
|
||||
<PackageReference Include="Quartz.Serialization.Json" Version="3.3.2" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\activities\Elsa.Activities.Email\Elsa.Activities.Email.csproj" />
|
||||
<ProjectReference Include="..\..\activities\Elsa.Activities.Temporal.Hangfire\Elsa.Activities.Temporal.Hangfire.csproj" />
|
||||
<ProjectReference Include="..\..\core\Elsa\Elsa.csproj" />
|
||||
<ProjectReference Include="..\..\persistence\Elsa.Persistence.EntityFramework\Elsa.Persistence.EntityFramework.Sqlite\Elsa.Persistence.EntityFramework.Sqlite.csproj" />
|
||||
<ProjectReference Include="..\..\samples\persistence\Elsa.Samples.Persistence.EntityFramework\Elsa.Samples.Persistence.EntityFramework.csproj" />
|
||||
<ProjectReference Include="..\Elsa.Server.Api\Elsa.Server.Api.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
3
src/server/Elsa.Server.Host/FodyWeavers.xml
Normal file
3
src/server/Elsa.Server.Host/FodyWeavers.xml
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd">
|
||||
<ConfigureAwait ContinueOnCapturedContext="false" />
|
||||
</Weavers>
|
||||
4
src/server/Elsa.Server.Host/Models/SampleInputModel.cs
Normal file
4
src/server/Elsa.Server.Host/Models/SampleInputModel.cs
Normal file
|
|
@ -0,0 +1,4 @@
|
|||
namespace Elsa.Samples.Server.Host.Models
|
||||
{
|
||||
public record SampleInputModel(string Name);
|
||||
}
|
||||
30
src/server/Elsa.Server.Host/Program.cs
Normal file
30
src/server/Elsa.Server.Host/Program.cs
Normal file
|
|
@ -0,0 +1,30 @@
|
|||
using Elsa.Server.Host;
|
||||
using Microsoft.AspNetCore.Hosting;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
|
||||
namespace Elsa.Samples.Server.Host
|
||||
{
|
||||
public class Program
|
||||
{
|
||||
public static void Main(string[] args)
|
||||
{
|
||||
CreateHostBuilder(args).Build().Run();
|
||||
}
|
||||
|
||||
public static IHostBuilder CreateHostBuilder(string[] args) =>
|
||||
Microsoft.Extensions.Hosting.Host.CreateDefaultBuilder(args)
|
||||
.ConfigureWebHostDefaults(webBuilder => webBuilder
|
||||
.UseStaticWebAssets()
|
||||
.UseStartup<Startup>())
|
||||
// .UseOrleans(siloBuilder => siloBuilder
|
||||
// .UseLocalhostClustering()
|
||||
// .Configure<ClusterOptions>(options =>
|
||||
// {
|
||||
// options.ClusterId = "localhost";
|
||||
// options.ServiceId = "elsa-workflows";
|
||||
// })
|
||||
// .ConfigureApplicationParts(parts => parts.AddApplicationPart(typeof(IWorkflowDefinitionGrain).Assembly).WithReferences())
|
||||
// .Configure<EndpointOptions>(options => options.AdvertisedIPAddress = IPAddress.Loopback))
|
||||
;
|
||||
}
|
||||
}
|
||||
13
src/server/Elsa.Server.Host/Properties/launchSettings.json
Normal file
13
src/server/Elsa.Server.Host/Properties/launchSettings.json
Normal file
|
|
@ -0,0 +1,13 @@
|
|||
{
|
||||
"profiles": {
|
||||
"Elsa.Server.Host": {
|
||||
"commandName": "Project",
|
||||
"launchBrowser": false,
|
||||
"applicationUrl": "https://localhost:11000",
|
||||
"launchUrl": "https://localhost:11000/swagger",
|
||||
"environmentVariables": {
|
||||
"ASPNETCORE_ENVIRONMENT": "Development"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
84
src/server/Elsa.Server.Host/Startup.cs
Normal file
84
src/server/Elsa.Server.Host/Startup.cs
Normal file
|
|
@ -0,0 +1,84 @@
|
|||
using System;
|
||||
using Elsa.Persistence.EntityFramework.Core.Extensions;
|
||||
using Elsa.Persistence.EntityFramework.Sqlite;
|
||||
using Elsa.Persistence.EntityFramework.SqlServer;
|
||||
using Hangfire;
|
||||
using Microsoft.AspNetCore.Builder;
|
||||
using Microsoft.AspNetCore.Hosting;
|
||||
using Microsoft.Extensions.Configuration;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
|
||||
namespace Elsa.Server.Host
|
||||
{
|
||||
public class Startup
|
||||
{
|
||||
public Startup(IWebHostEnvironment environment, IConfiguration configuration)
|
||||
{
|
||||
Environment = environment;
|
||||
Configuration = configuration;
|
||||
}
|
||||
|
||||
private IWebHostEnvironment Environment { get; }
|
||||
private IConfiguration Configuration { get; }
|
||||
|
||||
public void ConfigureServices(IServiceCollection services)
|
||||
{
|
||||
var elsaSection = Configuration.GetSection("Elsa");
|
||||
var sqlServerConnectionString = Configuration.GetConnectionString("SqlServer");
|
||||
|
||||
services.AddControllers();
|
||||
|
||||
services
|
||||
//.AddRedis(Configuration.GetConnectionString("Redis"))
|
||||
.AddElsa(elsa => elsa
|
||||
//.WithContainerName(Configuration["ContainerName"] ?? System.Environment.MachineName)
|
||||
.UseEntityFrameworkPersistence(ef => ef.UseSqlServer(sqlServerConnectionString))
|
||||
//.UseRabbitMq(Configuration.GetConnectionString("RabbitMq"))
|
||||
//.UseRebusCacheSignal()
|
||||
//.UseRedisCacheSignal()
|
||||
// .ConfigureDistributedLockProvider(options => options.UseProviderFactory(sp => name =>
|
||||
// {
|
||||
// var connection = sp.GetRequiredService<IConnectionMultiplexer>();
|
||||
// return new RedisDistributedLock(name, connection.GetDatabase());
|
||||
// }))
|
||||
.AddConsoleActivities()
|
||||
.AddHttpActivities(elsaSection.GetSection("Http").Bind)
|
||||
.AddEmailActivities(elsaSection.GetSection("Smtp").Bind)
|
||||
//.AddQuartzTemporalActivities()
|
||||
.AddHangfireTemporalActivities(hangfire => hangfire.UseSqlServerStorage(sqlServerConnectionString), (_, hangfireServer) => hangfireServer.SchedulePollingInterval = TimeSpan.FromSeconds(5))
|
||||
.AddJavaScriptActivities()
|
||||
.AddActivitiesFrom<Startup>()
|
||||
.AddWorkflowsFrom<Startup>()
|
||||
);
|
||||
|
||||
// Elsa API endpoints.
|
||||
services
|
||||
.AddElsaApiEndpoints()
|
||||
.AddElsaSwagger();
|
||||
|
||||
// Allow arbitrary client browser apps to access the API for demo purposes only.
|
||||
// In a production environment, make sure to allow only origins you trust.
|
||||
services.AddCors(cors => cors.AddDefaultPolicy(policy => policy.AllowAnyHeader().AllowAnyMethod().AllowAnyOrigin().WithExposedHeaders("Content-Disposition")));
|
||||
}
|
||||
|
||||
public void Configure(IApplicationBuilder app)
|
||||
{
|
||||
if (Environment.IsDevelopment())
|
||||
{
|
||||
app.UseDeveloperExceptionPage();
|
||||
app.UseSwagger();
|
||||
app.UseSwaggerUI(c => c.SwaggerEndpoint("/swagger/v1/swagger.json", "Elsa"));
|
||||
}
|
||||
|
||||
app
|
||||
.UseCors()
|
||||
.UseHttpActivities()
|
||||
.UseRouting()
|
||||
.UseEndpoints(endpoints =>
|
||||
{
|
||||
endpoints.MapControllers();
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
10
src/server/Elsa.Server.Host/appsettings.Development.json
Normal file
10
src/server/Elsa.Server.Host/appsettings.Development.json
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
{
|
||||
"Logging": {
|
||||
"LogLevel": {
|
||||
"Default": "Debug",
|
||||
"Orleans": "Warning",
|
||||
"System": "Information",
|
||||
"Microsoft": "Warning"
|
||||
}
|
||||
}
|
||||
}
|
||||
27
src/server/Elsa.Server.Host/appsettings.json
Normal file
27
src/server/Elsa.Server.Host/appsettings.json
Normal file
|
|
@ -0,0 +1,27 @@
|
|||
{
|
||||
"Logging": {
|
||||
"LogLevel": {
|
||||
"Default": "Warning",
|
||||
"Microsoft": "Warning",
|
||||
"Orleans": "Warning",
|
||||
"Microsoft.Hosting.Lifetime": "Information"
|
||||
}
|
||||
},
|
||||
"AllowedHosts": "*",
|
||||
"ConnectionStrings": {
|
||||
"RabbitMq": "amqp://localhost:5672",
|
||||
"Redis": "localhost:6379,abortConnect=false",
|
||||
"SqlServer": ""
|
||||
},
|
||||
"Elsa": {
|
||||
"Http": {
|
||||
"BaseUrl": "https://localhost:11000",
|
||||
"BasePath": "/workflows"
|
||||
},
|
||||
"Smtp": {
|
||||
"Host": "localhost",
|
||||
"Port": "2525",
|
||||
"DefaultSender": "noreply@acme.com"
|
||||
}
|
||||
}
|
||||
}
|
||||
30
src/server/Elsa.Server.Host/docker-compose.yaml
Normal file
30
src/server/Elsa.Server.Host/docker-compose.yaml
Normal file
|
|
@ -0,0 +1,30 @@
|
|||
version: '3.7'
|
||||
|
||||
services:
|
||||
|
||||
mongodb:
|
||||
image: mongo
|
||||
ports:
|
||||
- "27017:27017"
|
||||
volumes:
|
||||
- /var/lib/mongodb:/var/lib/mongodb
|
||||
|
||||
azureblobstorage:
|
||||
image: mcr.microsoft.com/azure-blob-storage
|
||||
|
||||
redis:
|
||||
image: redis
|
||||
ports:
|
||||
- "6379:6379"
|
||||
|
||||
rabbitmq:
|
||||
image: "rabbitmq:3-management"
|
||||
ports:
|
||||
- "15672:15672"
|
||||
- "5672:5672"
|
||||
|
||||
smtp4dev:
|
||||
image: rnwood/smtp4dev:linux-amd64-3.1.0-ci0856
|
||||
ports:
|
||||
- "3000:80"
|
||||
- "2525:25"
|
||||
Loading…
Reference in a new issue