Remove mediator-based workflow definition handling

Replaced mediator requests and commands with direct store interactions in the workflow management module. Simplified service dependencies, eliminating unnecessary handlers and requests, and streamlined workflow-related operations for better maintainability and performance.
This commit is contained in:
Sipke Schoorstra 2025-02-18 19:39:37 +01:00
parent dff6e14093
commit bc005a7586
No known key found for this signature in database
GPG key ID: 5C10502B28A4268F
13 changed files with 39 additions and 160 deletions

View file

@ -15,7 +15,6 @@ using Elsa.Http.Services;
using Elsa.Http.Tasks;
using Elsa.Http.UIHints;
using Elsa.Workflows;
using Elsa.Workflows.Management.Requests;
using FluentStorage;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.StaticFiles;

View file

@ -1,14 +1,13 @@
using Elsa.Abstractions;
using Elsa.Common.Models;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Requests;
using Elsa.Workflows.Management;
using Elsa.Workflows.Models;
using JetBrains.Annotations;
namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.GetByDefinitionId;
[PublicAPI]
internal class GetByDefinitionId(IMediator mediator, IWorkflowDefinitionLinker linker) : ElsaEndpoint<Request>
internal class GetByDefinitionId(IWorkflowDefinitionStore store, IWorkflowDefinitionLinker linker) : ElsaEndpoint<Request>
{
public override void Configure()
{
@ -19,9 +18,8 @@ internal class GetByDefinitionId(IMediator mediator, IWorkflowDefinitionLinker l
public override async Task HandleAsync(Request request, CancellationToken cancellationToken)
{
var versionOptions = request.VersionOptions != null ? VersionOptions.FromString(request.VersionOptions) : VersionOptions.Latest;
var handle = WorkflowDefinitionHandle.ByDefinitionId(request.DefinitionId, versionOptions);
var findRequest = new FindWorkflowDefinitionRequest(handle);
var definition = await mediator.SendAsync(findRequest, cancellationToken);
var filter = WorkflowDefinitionHandle.ByDefinitionId(request.DefinitionId, versionOptions).ToFilter();
var definition = await store.FindAsync(filter, cancellationToken);
if (definition == null)
{

View file

@ -1,6 +0,0 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Management.Commands;
public record SaveWorkflowDefinitionCommand(WorkflowDefinition WorkflowDefinition) : ICommand;

View file

@ -1,8 +1,5 @@
using Elsa.Features.Abstractions;
using Elsa.Features.Services;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Handlers.Commands;
using Elsa.Workflows.Management.Handlers.Requests;
using Elsa.Workflows.Management.Stores;
using Microsoft.Extensions.DependencyInjection;
@ -17,20 +14,11 @@ public class WorkflowDefinitionsFeature(IModule module) : FeatureBase(module)
/// The factory to create new instances of <see cref="IWorkflowDefinitionStore"/>.
/// </summary>
public Func<IServiceProvider, IWorkflowDefinitionStore> WorkflowDefinitionStore { get; set; } = sp => sp.GetRequiredService<MemoryWorkflowDefinitionStore>();
public Type FindWorkflowDefinitionHandler { get; set; } = typeof(FindWorkflowDefinitionHandler);
public Type FindLastVersionOfWorkflowDefinitionHandler { get; set; } = typeof(FindLastVersionOfWorkflowDefinitionHandler);
public Type FindLatestOrPublishedWorkflowDefinitionsHandler { get; set; } = typeof(FindLatestOrPublishedWorkflowDefinitionsHandler);
public Type SaveWorkflowDefinitionHandler { get; set; } = typeof(SaveWorkflowDefinitionHandler);
/// <inheritdoc />
public override void Apply()
{
Services
.AddScoped(WorkflowDefinitionStore)
.AddScoped(typeof(IRequestHandler), FindWorkflowDefinitionHandler)
.AddScoped(typeof(IRequestHandler), FindLastVersionOfWorkflowDefinitionHandler)
.AddScoped(typeof(IRequestHandler), FindLatestOrPublishedWorkflowDefinitionsHandler)
.AddScoped(typeof(ICommandHandler), SaveWorkflowDefinitionHandler)
;
.AddScoped(WorkflowDefinitionStore);
}
}

View file

@ -1,14 +0,0 @@
using Elsa.Mediator.Contracts;
using Elsa.Mediator.Models;
using Elsa.Workflows.Management.Commands;
namespace Elsa.Workflows.Management.Handlers.Commands;
public class SaveWorkflowDefinitionHandler(IWorkflowDefinitionStore store) : ICommandHandler<SaveWorkflowDefinitionCommand>
{
public async Task<Unit> HandleAsync(SaveWorkflowDefinitionCommand command, CancellationToken cancellationToken)
{
await store.SaveAsync(command.WorkflowDefinition, cancellationToken);
return Unit.Instance;
}
}

View file

@ -1,18 +0,0 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Management.Requests;
namespace Elsa.Workflows.Management.Handlers.Requests;
public class FindLastVersionOfWorkflowDefinitionHandler(IWorkflowDefinitionStore store) : IRequestHandler<FindLastVersionOfWorkflowDefinitionRequest, WorkflowDefinition?>
{
public Task<WorkflowDefinition?> HandleAsync(FindLastVersionOfWorkflowDefinitionRequest request, CancellationToken cancellationToken)
{
var filter = new WorkflowDefinitionFilter()
{
DefinitionId = request.DefinitionId,
};
return store.FindLastVersionAsync(filter, cancellationToken: cancellationToken);
}
}

View file

@ -1,31 +0,0 @@
using Elsa.Common.Entities;
using Elsa.Common.Models;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Management.Requests;
namespace Elsa.Workflows.Management.Handlers.Requests;
public class FindLatestOrPublishedWorkflowDefinitionsHandler(IWorkflowDefinitionStore store) : IRequestHandler<FindLatestOrPublishedWorkflowDefinitionsRequest, ICollection<WorkflowDefinition>>
{
public Task<WorkflowDefinition?> HandleAsync(FindLastVersionOfWorkflowDefinitionRequest request, CancellationToken cancellationToken)
{
var filter = new WorkflowDefinitionFilter()
{
DefinitionId = request.DefinitionId,
};
return store.FindLastVersionAsync(filter, cancellationToken: cancellationToken);
}
public async Task<ICollection<WorkflowDefinition>> HandleAsync(FindLatestOrPublishedWorkflowDefinitionsRequest request, CancellationToken cancellationToken)
{
var filter = new WorkflowDefinitionFilter
{
DefinitionId = request.DefinitionId,
VersionOptions = VersionOptions.LatestOrPublished
};
var order = new WorkflowDefinitionOrder<int>(x => x.Version, OrderDirection.Descending);
return (await store.FindManyAsync(filter, order, cancellationToken)).ToList();
}
}

View file

@ -1,18 +0,0 @@
using Elsa.Common.Entities;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Management.Requests;
namespace Elsa.Workflows.Management.Handlers.Requests;
public class FindWorkflowDefinitionHandler(IWorkflowDefinitionStore store) : IRequestHandler<FindWorkflowDefinitionRequest, WorkflowDefinition?>
{
public async Task<WorkflowDefinition?> HandleAsync(FindWorkflowDefinitionRequest request, CancellationToken cancellationToken)
{
var filter = request.Handle.ToFilter();
var order = new WorkflowDefinitionOrder<int>(x => x.Version, OrderDirection.Descending);
var definition = (await store.FindManyAsync(filter, order, cancellationToken: cancellationToken)).FirstOrDefault();
return definition;
}
}

View file

@ -1,10 +0,0 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Management.Requests;
/// <summary>
/// A request to find the last version of a workflow definition.
/// </summary>
/// <param name="DefinitionId">The ID of the workflow definition.</param>
public record FindLastVersionOfWorkflowDefinitionRequest(string DefinitionId) : IRequest<WorkflowDefinition?>;

View file

@ -1,9 +0,0 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Management.Requests;
/// <summary>
/// A request to find the latest or published workflow definitions.
/// </summary>
public record FindLatestOrPublishedWorkflowDefinitionsRequest(string DefinitionId) : IRequest<ICollection<WorkflowDefinition>>;

View file

@ -1,12 +0,0 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Requests;
/// <summary>
/// A request to find a workflow definition.
/// </summary>
/// <param name="DefinitionId">The ID of the workflow definition.</param>
/// <param name="VersionOptions">The version options.</param>
public record FindWorkflowDefinitionRequest(WorkflowDefinitionHandle Handle) : IRequest<WorkflowDefinition?>;

View file

@ -1,13 +1,13 @@
using Elsa.Common;
using Elsa.Common.Entities;
using Elsa.Common.Models;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Management.Commands;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Management.Materializers;
using Elsa.Workflows.Management.Models;
using Elsa.Workflows.Management.Notifications;
using Elsa.Workflows.Management.Requests;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Services;
@ -69,8 +69,8 @@ public class WorkflowDefinitionPublisher(
/// <inheritdoc />
public async Task<PublishWorkflowDefinitionResult> PublishAsync(string definitionId, CancellationToken cancellationToken = default)
{
var handle = WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Latest);
var definition = await mediator.SendAsync(new FindWorkflowDefinitionRequest(handle), cancellationToken);
var filter = WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Latest).ToFilter();
var definition = await workflowDefinitionStore.FindAsync(filter, cancellationToken);
if (definition == null)
return new(false, new List<WorkflowValidationError>
@ -94,14 +94,18 @@ public class WorkflowDefinitionPublisher(
var definitionId = definition.DefinitionId;
// Reset current latest and published definitions.
var publishedWorkflows = await mediator.SendAsync(new FindLatestOrPublishedWorkflowDefinitionsRequest(definitionId), cancellationToken);
var publishedWorkflows = await workflowDefinitionStore.FindManyAsync(new()
{
DefinitionId = definitionId,
VersionOptions = VersionOptions.LatestOrPublished
}, cancellationToken);
foreach (var publishedAndOrLatestWorkflow in publishedWorkflows)
{
var isPublished = publishedAndOrLatestWorkflow.IsPublished;
publishedAndOrLatestWorkflow.IsPublished = false;
publishedAndOrLatestWorkflow.IsLatest = false;
await mediator.SendAsync(new SaveWorkflowDefinitionCommand(publishedAndOrLatestWorkflow), cancellationToken);
await workflowDefinitionStore.SaveAsync(publishedAndOrLatestWorkflow, cancellationToken);
if (isPublished)
await mediator.SendAsync(new WorkflowDefinitionVersionRetracted(publishedAndOrLatestWorkflow), cancellationToken);
@ -110,7 +114,7 @@ public class WorkflowDefinitionPublisher(
// Save the newly published definition.
definition.IsPublished = true;
definition = Initialize(definition);
await mediator.SendAsync(new SaveWorkflowDefinitionCommand(definition), cancellationToken);
await workflowDefinitionStore.SaveAsync(definition, cancellationToken);
var affectedWorkflows = new AffectedWorkflows(new List<WorkflowDefinition>());
await mediator.SendAsync(new WorkflowDefinitionPublished(definition, affectedWorkflows), cancellationToken);
@ -120,8 +124,8 @@ public class WorkflowDefinitionPublisher(
/// <inheritdoc />
public async Task<WorkflowDefinition?> RetractAsync(string definitionId, CancellationToken cancellationToken = default)
{
var handle = WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Published);
var definition = await mediator.SendAsync(new FindWorkflowDefinitionRequest(handle), cancellationToken);
var filter = WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Published).ToFilter();
var definition = await workflowDefinitionStore.FindAsync(filter, cancellationToken);
if (definition == null)
return null;
@ -138,7 +142,7 @@ public class WorkflowDefinitionPublisher(
definition.IsPublished = false;
await mediator.SendAsync(new WorkflowDefinitionRetracting(definition), cancellationToken);
await mediator.SendAsync(new SaveWorkflowDefinitionCommand(definition), cancellationToken);
await workflowDefinitionStore.SaveAsync(definition, cancellationToken);
await mediator.SendAsync(new WorkflowDefinitionRetracted(definition), cancellationToken);
return definition;
}
@ -146,11 +150,17 @@ public class WorkflowDefinitionPublisher(
/// <inheritdoc />
public async Task<WorkflowDefinition?> GetDraftAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default)
{
var findLastVersionRequest = new FindLastVersionOfWorkflowDefinitionRequest(definitionId);
var lastVersion = await mediator.SendAsync(findLastVersionRequest, cancellationToken);
var handle = WorkflowDefinitionHandle.ByDefinitionId(definitionId, versionOptions);
var findRequest = new FindWorkflowDefinitionRequest(handle);
var definition = await mediator.SendAsync(findRequest, cancellationToken) ?? lastVersion;
var filter = new WorkflowDefinitionFilter
{
DefinitionId = definitionId,
VersionOptions = versionOptions
};
var order = new WorkflowDefinitionOrder<int>(x => x.Version, OrderDirection.Descending);
var lastVersion = await workflowDefinitionStore.FindLastVersionAsync(new()
{
DefinitionId = definitionId
}, cancellationToken);
var definition = await workflowDefinitionStore.FindAsync(filter, order, cancellationToken) ?? lastVersion;
if (definition == null!)
return null;
@ -174,13 +184,17 @@ public class WorkflowDefinitionPublisher(
{
var draft = definition;
var definitionId = definition.DefinitionId;
var lastVersion = await mediator.SendAsync(new FindLastVersionOfWorkflowDefinitionRequest(definitionId), cancellationToken);
var filter = new WorkflowDefinitionFilter
{
DefinitionId = definitionId
};
var lastVersion = await workflowDefinitionStore.FindLastVersionAsync(filter, cancellationToken);
draft.Version = draft.Id == lastVersion?.Id ? lastVersion.Version : lastVersion?.Version + 1 ?? 1;
draft.IsLatest = true;
draft = Initialize(draft);
await mediator.SendAsync(new SaveWorkflowDefinitionCommand(draft), cancellationToken);
await workflowDefinitionStore.SaveAsync(draft, cancellationToken);
if (lastVersion is null)
await mediator.SendAsync(new WorkflowDefinitionCreated(definition), cancellationToken);
@ -188,7 +202,7 @@ public class WorkflowDefinitionPublisher(
if (lastVersion is { IsPublished: true, IsLatest: true })
{
lastVersion.IsLatest = false;
await mediator.SendAsync(new SaveWorkflowDefinitionCommand(lastVersion), cancellationToken);
await workflowDefinitionStore.SaveAsync(lastVersion, cancellationToken);
}
return draft;

View file

@ -1,8 +1,6 @@
using Elsa.Common.Models;
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Management.Requests;
using Elsa.Workflows.Models;
namespace Elsa.Workflows.Management.Services;
@ -11,7 +9,6 @@ namespace Elsa.Workflows.Management.Services;
public class WorkflowDefinitionService(
IWorkflowDefinitionStore workflowDefinitionStore,
IWorkflowGraphBuilder workflowGraphBuilder,
IMediator mediator,
Func<IEnumerable<IWorkflowMaterializer>> materializers)
: IWorkflowDefinitionService
{
@ -45,7 +42,8 @@ public class WorkflowDefinitionService(
/// <inheritdoc />
public async Task<WorkflowDefinition?> FindWorkflowDefinitionAsync(WorkflowDefinitionHandle handle, CancellationToken cancellationToken = default)
{
return await mediator.SendAsync(new FindWorkflowDefinitionRequest(handle), cancellationToken);
var filter = handle.ToFilter();
return await workflowDefinitionStore.FindAsync(filter, cancellationToken);
}
/// <inheritdoc />