diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkDispatch/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkDispatch/Endpoint.cs new file mode 100644 index 000000000..df1b7d3d9 --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkDispatch/Endpoint.cs @@ -0,0 +1,73 @@ +using Elsa.Abstractions; +using Elsa.Common.Models; +using Elsa.Workflows.Core.Contracts; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Requests; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.BulkDispatch; + +[PublicAPI] +internal class Endpoint : ElsaEndpoint +{ + private readonly IWorkflowDefinitionStore _store; + private readonly IWorkflowDispatcher _workflowDispatcher; + private readonly IIdentityGenerator _identityGenerator; + + public Endpoint(IWorkflowDefinitionStore store, IWorkflowDispatcher workflowDispatcher, IIdentityGenerator identityGenerator) + { + _store = store; + _workflowDispatcher = workflowDispatcher; + _identityGenerator = identityGenerator; + } + + public override void Configure() + { + Post("/workflow-definitions/{definitionId}/bulk-dispatch"); + ConfigurePermissions("exec:workflow-definitions"); + } + + public override async Task HandleAsync(Request request, CancellationToken cancellationToken) + { + var definitionId = request.DefinitionId; + var versionOptions = request.VersionOptions ?? VersionOptions.Published; + + var exists = await _store.AnyAsync( + new WorkflowDefinitionFilter + { + DefinitionId = definitionId, + VersionOptions = versionOptions + }, + cancellationToken); + + if (!exists) + { + await SendNotFoundAsync(cancellationToken); + return; + } + + var instanceIds = new List(); + + for (var i = 0; i < request.Count; i++) + { + var instanceId = _identityGenerator.GenerateId(); + var triggerActivityId = request.TriggerActivityId; + var input = (IDictionary?)request.Input; + var dispatchRequest = new DispatchWorkflowDefinitionRequest + { + DefinitionId = definitionId, + VersionOptions = versionOptions, + Input = input, + InstanceId = instanceId, + TriggerActivityId = triggerActivityId + }; + + await _workflowDispatcher.DispatchAsync(dispatchRequest, cancellationToken); + instanceIds.Add(instanceId); + } + + await SendOkAsync(new Response(instanceIds), cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkDispatch/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkDispatch/Models.cs new file mode 100644 index 000000000..a3c4b7dea --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/BulkDispatch/Models.cs @@ -0,0 +1,20 @@ +using System.Text.Json.Serialization; +using Elsa.Common.Models; +using Elsa.Workflows.Core.Serialization.Converters; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.BulkDispatch; + +internal class Request +{ + public string DefinitionId { get; set; } = default!; + public string? TriggerActivityId { get; set; } + + public VersionOptions? VersionOptions { get; set; } + + [JsonConverter(typeof(ExpandoObjectConverterFactory))] + public object? Input { get; set; } + + public int Count { get; set; } = 1; +} + +internal record Response(ICollection WorkflowInstanceIds); \ No newline at end of file