Add alteration notifications and management (#5097)

* Add alteration notifications and management

Implement notification publishing for completed alteration plans and add related handlers and management logic to coordinate alterations workflow triggering and completion status updates.

* Refactor queue address configuration in MassTransit

The initialization of the queue address and mapping for the `RunAlterationJob` has been moved out of the `Apply` method to the module configuration.

* Refactor CompleteAlterationPlan activity description

The activity summary comment in CompleteAlterationPlan.cs was updated to accurately describe its functionality. The description has been changed from "Submits an alteration plan for execution" to "Marks an alteration plan as completed."
This commit is contained in:
Sipke Schoorstra 2024-03-20 10:58:59 +01:00 committed by GitHub
parent 05fd1440c1
commit 8079dd0de1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
13 changed files with 174 additions and 70 deletions

View file

@ -0,0 +1,10 @@
using Elsa.Alterations.Core.Entities;
namespace Elsa.Alterations.Core.Contracts;
public interface IAlterationPlanManager
{
Task<AlterationPlan?> GetPlanAsync(string planId, CancellationToken cancellationToken = default);
Task<bool> GetIsAllJobsCompletedAsync(string planId, CancellationToken cancellationToken = default);
Task CompletePlanAsync(AlterationPlan plan, CancellationToken cancellationToken = default);
}

View file

@ -0,0 +1,9 @@
using Elsa.Alterations.Core.Entities;
using Elsa.Mediator.Contracts;
namespace Elsa.Alterations.Core.Notifications;
/// <summary>
/// A notification that is published when an alteration plan is completed.
/// </summary>
public record AlterationPlanCompleted(AlterationPlan Plan) : INotification;

View file

@ -28,15 +28,15 @@ public class MassTransitAlterationsFeature : FeatureBase
public override void Configure()
{
Module.Configure<AlterationsFeature>(feature => feature.AlterationJobDispatcherFactory = sp => sp.GetRequiredService<MassTransitAlterationJobDispatcher>());
Module.AddMassTransitConsumer<RunAlterationJobConsumer>();
var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer<RunAlterationJobConsumer>();
var queueAddress = new Uri($"queue:elsa-{queueName}");
EndpointConvention.Map<RunAlterationJob>(queueAddress);
Module.AddMassTransitConsumer<RunAlterationJobConsumer>(queueName);
}
/// <inheritdoc />
public override void Apply()
{
var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer<RunAlterationJobConsumer>();
var queueAddress = new Uri($"queue:elsa-{queueName}");
EndpointConvention.Map<RunAlterationJob>(queueAddress);
Services.AddScoped<MassTransitAlterationJobDispatcher>();
}
}

View file

@ -0,0 +1,48 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using Elsa.Alterations.Core.Contracts;
using Elsa.Workflows;
using Elsa.Workflows.Attributes;
using Elsa.Workflows.Exceptions;
using Elsa.Workflows.Memory;
using Elsa.Workflows.Models;
namespace Elsa.Alterations.Activities;
/// <summary>
/// Marks an alteration plan as completed.
/// </summary>
[Browsable(false)]
[Activity("Elsa", "Alterations", "Dispatches jobs for the specified Alteration Plan", Kind = ActivityKind.Job)]
public class CompleteAlterationPlan : CodeActivity
{
/// <inheritdoc />
public CompleteAlterationPlan(Variable<string> planId, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
PlanId = new Input<string>(planId);
}
/// <inheritdoc />
public CompleteAlterationPlan([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
}
/// <summary>
/// The ID of the alteration plan.
/// </summary>
public Input<string> PlanId { get; set; } = default!;
/// <inheritdoc />
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
{
var cancellationToken = context.CancellationToken;
var planId = context.Get(PlanId)!;
var manager = context.GetRequiredService<IAlterationPlanManager>();
var plan = await manager.GetPlanAsync(planId, cancellationToken);
if (plan == null)
throw new FaultException($"Alteration Plan with ID {planId} not found.");
await manager.CompletePlanAsync(plan, cancellationToken);
}
}

View file

@ -1,20 +1,13 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using Elsa.Alterations.Core.Contracts;
using Elsa.Alterations.Core.Entities;
using Elsa.Alterations.Core.Enums;
using Elsa.Alterations.Core.Filters;
using Elsa.Common.Contracts;
using Elsa.Workflows;
using Elsa.Workflows.Attributes;
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Exceptions;
using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Memory;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Filters;
namespace Elsa.Alterations.Activities;
@ -72,5 +65,9 @@ public class DispatchAlterationJobs : CodeActivity
var alterationJobDispatcher = context.GetRequiredService<IAlterationJobDispatcher>();
foreach (var jobId in alterationJobIds)
await alterationJobDispatcher.DispatchAsync(jobId, cancellationToken);
// Update status.
plan.Status = AlterationPlanStatus.Running;
await alterationPlanStore.SaveAsync(plan, cancellationToken);
}
}

View file

@ -6,6 +6,7 @@ using Elsa.Alterations.Core.Enums;
using Elsa.Alterations.Core.Filters;
using Elsa.Alterations.Core.Models;
using Elsa.Common.Contracts;
using Elsa.Extensions;
using Elsa.Workflows;
using Elsa.Workflows.Attributes;
using Elsa.Workflows.Contracts;
@ -24,7 +25,7 @@ namespace Elsa.Alterations.Activities;
/// </summary>
[Browsable(false)]
[Activity("Elsa", "Alterations", "Generates jobs for the specified Alteration Plan", Kind = ActivityKind.Job)]
public class GenerateAlterationJobs : CodeActivity
public class GenerateAlterationJobs : CodeActivity<int>
{
/// <inheritdoc />
public GenerateAlterationJobs([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
@ -47,10 +48,14 @@ public class GenerateAlterationJobs : CodeActivity
{
var plan = await GetPlanAsync(context);
await UpdatePlanStatusAsync(context, plan);
var workflowInstanceIds = await FindMatchingWorkflowInstanceIdsAsync(context, plan.WorkflowInstanceFilter);
await GenerateJobsAsync(context, plan, workflowInstanceIds);
var workflowInstanceIds = (await FindMatchingWorkflowInstanceIdsAsync(context, plan.WorkflowInstanceFilter)).ToList();
if (workflowInstanceIds.Any())
await GenerateJobsAsync(context, plan, workflowInstanceIds);
context.SetResult(workflowInstanceIds.Count);
}
private async Task<AlterationPlan> GetPlanAsync(ActivityExecutionContext context)
{
var cancellationToken = context.CancellationToken;
@ -64,7 +69,7 @@ public class GenerateAlterationJobs : CodeActivity
if (plan == null)
throw new FaultException($"Alteration Plan with ID {planId} not found.");
return plan;
}
@ -108,10 +113,10 @@ public class GenerateAlterationJobs : CodeActivity
workflowInstanceIds.UnionWith(matchingWorkflowInstanceIds);
}
}
return workflowInstanceIds;
}
private async Task GenerateJobsAsync(ActivityExecutionContext context, AlterationPlan plan, IEnumerable<string> workflowInstanceIds)
{
var cancellationToken = context.CancellationToken;

View file

@ -1,7 +1,6 @@
using Elsa.Alterations.Core.Contracts;
using Elsa.Alterations.Core.Entities;
using Elsa.Alterations.Core.Extensions;
using Elsa.Alterations.Core.Services;
using Elsa.Alterations.Core.Stores;
using Elsa.Alterations.Extensions;
using Elsa.Alterations.Services;
@ -60,6 +59,7 @@ public class AlterationsFeature : FeatureBase
/// <inheritdoc />
public override void Apply()
{
Services.AddScoped<IAlterationPlanManager, AlterationPlanManager>();
Services.AddAlterations();
Services.AddAlterationsCore();
Services.AddScoped<BackgroundAlterationJobDispatcher>();

View file

@ -1,14 +1,6 @@
using Elsa.Alterations.Activities;
using Elsa.Alterations.Bookmarks;
using Elsa.Alterations.Core.Contracts;
using Elsa.Alterations.Core.Enums;
using Elsa.Alterations.Core.Filters;
using Elsa.Alterations.Core.Notifications;
using Elsa.Common.Contracts;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Helpers;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Requests;
using JetBrains.Annotations;
namespace Elsa.Alterations.Handlers;
@ -17,53 +9,19 @@ namespace Elsa.Alterations.Handlers;
/// Handles <see cref="AlterationJobCompleted"/> notifications and updates the plan status if all jobs are completed.
/// </summary>
[UsedImplicitly]
public class AlterationJobCompletedHandler : INotificationHandler<AlterationJobCompleted>
public class AlterationJobCompletedHandler(IAlterationPlanManager manager) : INotificationHandler<AlterationJobCompleted>
{
private readonly IWorkflowDispatcher _workflowDispatcher;
private readonly IAlterationPlanStore _alterationPlanStore;
private readonly IAlterationJobStore _alterationJobStore;
private readonly ISystemClock _systemClock;
/// <summary>
/// Initializes a new instance of the <see cref="AlterationJobCompletedHandler"/> class.
/// </summary>
public AlterationJobCompletedHandler(IWorkflowDispatcher workflowDispatcher, IAlterationPlanStore alterationPlanStore, IAlterationJobStore alterationJobStore, ISystemClock systemClock)
{
_workflowDispatcher = workflowDispatcher;
_alterationPlanStore = alterationPlanStore;
_alterationJobStore = alterationJobStore;
_systemClock = systemClock;
}
/// <inheritdoc />
public async Task HandleAsync(AlterationJobCompleted notification, CancellationToken cancellationToken)
{
var job = notification.Job;
var planId = job.PlanId;
var planFilter = new AlterationPlanFilter { Id = planId };
var plan = (await _alterationPlanStore.FindAsync(planFilter, cancellationToken))!;
// Check if all jobs are completed.
var jobFilter = new AlterationJobFilter
{
PlanId = planId,
Statuses = new[] { AlterationJobStatus.Pending, AlterationJobStatus.Running }
};
var allJobsCompleted = await _alterationJobStore.CountAsync(jobFilter, cancellationToken) == 0;
var plan = (await manager.GetPlanAsync(planId, cancellationToken))!;
var allJobsCompleted = await manager.GetIsAllJobsCompletedAsync(planId, cancellationToken);
if(!allJobsCompleted)
return;
// Update plan status.
plan.Status = AlterationPlanStatus.Completed;
plan.CompletedAt = _systemClock.UtcNow;
await _alterationPlanStore.SaveAsync(plan, cancellationToken);
// Trigger any workflow instances that are waiting for the plan to complete.
var bookmarkPayload = new AlterationPlanCompletedPayload(planId);
var triggerRequest = new DispatchTriggerWorkflowsRequest(ActivityTypeNameHelper.GenerateTypeName<AlterationPlanCompleted>(), bookmarkPayload);
await _workflowDispatcher.DispatchAsync(triggerRequest, cancellationToken);
await manager.CompletePlanAsync(plan, cancellationToken);
}
}

View file

@ -0,0 +1,26 @@
using Elsa.Alterations.Bookmarks;
using Elsa.Alterations.Core.Notifications;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Helpers;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Requests;
using JetBrains.Annotations;
namespace Elsa.Alterations.Handlers;
/// <summary>
/// Handles <see cref="AlterationPlanCompleted"/> notifications and triggers any workflows that are waiting for the plan to complete.
/// </summary>
[UsedImplicitly]
public class AlterationPlanCompletedHandler(IWorkflowDispatcher workflowDispatcher) : INotificationHandler<AlterationPlanCompleted>
{
/// <inheritdoc />
public async Task HandleAsync(AlterationPlanCompleted notification, CancellationToken cancellationToken)
{
// Trigger any workflow instances that are waiting for the plan to complete.
var planId = notification.Plan.Id;
var bookmarkPayload = new AlterationPlanCompletedPayload(planId);
var triggerRequest = new DispatchTriggerWorkflowsRequest(ActivityTypeNameHelper.GenerateTypeName<Activities.AlterationPlanCompleted>(), bookmarkPayload);
await workflowDispatcher.DispatchAsync(triggerRequest, cancellationToken);
}
}

View file

@ -0,0 +1,43 @@
using Elsa.Alterations.Core.Contracts;
using Elsa.Alterations.Core.Entities;
using Elsa.Alterations.Core.Enums;
using Elsa.Alterations.Core.Filters;
using Elsa.Alterations.Core.Notifications;
using Elsa.Common.Contracts;
using Elsa.Mediator.Contracts;
namespace Elsa.Alterations.Services;
/// <inheritdoc />
public class AlterationPlanManager(IAlterationPlanStore planStore, IAlterationJobStore jobStore, ISystemClock systemClock, INotificationSender notificationSender) : IAlterationPlanManager
{
/// <inheritdoc />
public async Task<AlterationPlan?> GetPlanAsync(string planId, CancellationToken cancellationToken = default)
{
var planFilter = new AlterationPlanFilter { Id = planId };
return await planStore.FindAsync(planFilter, cancellationToken);
}
/// <inheritdoc />
public async Task<bool> GetIsAllJobsCompletedAsync(string planId, CancellationToken cancellationToken = default)
{
// Check if all jobs are completed.
var jobFilter = new AlterationJobFilter
{
PlanId = planId,
Statuses = new[] { AlterationJobStatus.Pending, AlterationJobStatus.Running }
};
return await jobStore.CountAsync(jobFilter, cancellationToken) == 0;
}
/// <inheritdoc />
public async Task CompletePlanAsync(AlterationPlan plan, CancellationToken cancellationToken = default)
{
plan.Status = AlterationPlanStatus.Completed;
plan.CompletedAt = systemClock.UtcNow;
await planStore.SaveAsync(plan, cancellationToken);
await notificationSender.SendAsync(new AlterationPlanCompleted(plan), cancellationToken);
}
}

View file

@ -6,7 +6,7 @@ using Elsa.Alterations.Core.Notifications;
using Elsa.Common.Contracts;
using Elsa.Mediator.Contracts;
namespace Elsa.Alterations.Core.Services;
namespace Elsa.Alterations.Services;
/// <inheritdoc />
public class DefaultAlterationJobRunner : IAlterationJobRunner

View file

@ -9,7 +9,7 @@ using Elsa.Workflows.Management.Contracts;
using Elsa.Workflows.Runtime.Contracts;
using Microsoft.Extensions.Logging;
namespace Elsa.Alterations.Core.Services;
namespace Elsa.Alterations.Services;
/// <inheritdoc />
public class DefaultAlterationRunner : IAlterationRunner

View file

@ -13,7 +13,7 @@ namespace Elsa.Alterations.Workflows;
public class ExecuteAlterationPlanWorkflow : WorkflowBase
{
internal const string WorkflowDefinitionId = "Elsa.Alterations.ExecuteAlterationPlan";
/// <inheritdoc />
protected override void Build(IWorkflowBuilder builder)
{
@ -21,6 +21,7 @@ public class ExecuteAlterationPlanWorkflow : WorkflowBase
builder.AsSystemWorkflow();
var plan = builder.WithInput<AlterationPlanParams>("Plan", "The parameters for the new plan");
var planId = builder.WithVariable<string>();
var jobCount = builder.WithVariable<int>();
builder.Root = new Sequence
{
@ -32,8 +33,15 @@ public class ExecuteAlterationPlanWorkflow : WorkflowBase
Result = new(planId)
},
new Correlate(planId),
new GenerateAlterationJobs(planId),
new DispatchAlterationJobs(planId),
new GenerateAlterationJobs(planId)
{
Result = new(jobCount)
},
new If(context => jobCount.Get(context) > 0)
{
Then = new DispatchAlterationJobs(planId),
Else = new CompleteAlterationPlan(planId)
},
new AlterationPlanCompleted(planId),
}
};