diff --git a/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs b/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs index e0ffd9172..35195a271 100644 --- a/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs +++ b/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs @@ -2,6 +2,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Signaling.Services; using Microsoft.AspNetCore.Mvc; +using Open.Linq.AsyncExtensions; namespace Elsa.Activities.Http.Endpoints.Signals { @@ -20,10 +21,8 @@ namespace Elsa.Activities.Http.Endpoints.Signals [HttpGet, HttpPost] public async Task Handle(string token, CancellationToken cancellationToken) { - if (!await _signaler.DispatchSignalTokenAsync(token, cancellationToken: cancellationToken)) - return NotFound(); - - return Accepted(); + var pendingWorkflows = await _signaler.DispatchSignalTokenAsync(token, cancellationToken: cancellationToken).ToList(); + return Accepted(pendingWorkflows); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs b/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs index 50022d91c..a24bd7dc9 100644 --- a/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs +++ b/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs @@ -2,6 +2,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Signaling.Services; using Microsoft.AspNetCore.Mvc; +using Open.Linq.AsyncExtensions; namespace Elsa.Activities.Http.Endpoints.Signals { @@ -16,12 +17,11 @@ namespace Elsa.Activities.Http.Endpoints.Signals [HttpGet, HttpPost] public async Task Handle(string token, CancellationToken cancellationToken) { - if (!await _signaler.TriggerSignalTokenAsync(token, cancellationToken: cancellationToken)) - return NotFound(); + var result = await _signaler.TriggerSignalTokenAsync(token, cancellationToken: cancellationToken).ToList(); return HttpContext.Response.HasStarted ? new EmptyResult() - : Accepted(); + : Ok(result); } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowLaunchpad.cs b/src/core/Elsa.Abstractions/Services/IWorkflowLaunchpad.cs index 5f32596f6..c11dae68d 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowLaunchpad.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowLaunchpad.cs @@ -87,7 +87,7 @@ namespace Elsa.Services /// /// Collects and executes workflows that are ready for execution. This takes into account both resumable (suspended) workflows as well as startable workflows. /// - Task> CollectAndExecuteWorkflowsAsync(CollectWorkflowsContext context, object? input = default, CancellationToken cancellationToken = default); + Task> CollectAndExecuteWorkflowsAsync(CollectWorkflowsContext context, object? input = null, CancellationToken cancellationToken = default); /// /// Collects and dispatches workflows that are ready for execution. This takes into account both resumable (suspended) workflows as well as startable workflows. diff --git a/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs b/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs index f13f2bf6e..e64cfe6fb 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs @@ -1,5 +1,7 @@ -using System.Threading; +using System.Collections.Generic; +using System.Threading; using System.Threading.Tasks; +using Elsa.Services.Models; namespace Elsa.Activities.Signaling.Services { @@ -8,21 +10,21 @@ namespace Elsa.Activities.Signaling.Services /// /// Runs all workflows that start with or are blocked on the activity. /// - Task TriggerSignalTokenAsync(string signalToken, object? input = default, CancellationToken cancellationToken = default); + Task> TriggerSignalTokenAsync(string signalToken, object? input = default, CancellationToken cancellationToken = default); /// /// Runs all workflows that start with or are blocked on the activity. /// - Task TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default); + Task> TriggerSignalAsync(string signal, object? input = null, string? workflowInstanceId = null, string? correlationId = null, CancellationToken cancellationToken = default); /// /// Dispatches all workflows that start with or are blocked on the activity. /// - Task DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default); + Task> DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default); /// /// Dispatches all workflows that start with or are blocked on the activity. /// - Task DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default); + Task> DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs b/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs index ca56abd08..c93dfa3fe 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs @@ -1,9 +1,10 @@ -using System.Threading; +using System.Collections.Generic; +using System.Linq; +using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Signaling.Models; -using Elsa.Dispatch; using Elsa.Services; -using MediatR; +using Elsa.Services.Models; namespace Elsa.Activities.Signaling.Services { @@ -21,18 +22,17 @@ namespace Elsa.Activities.Signaling.Services _tokenService = tokenService; } - public async Task TriggerSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default) + public async Task> TriggerSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default) { if (!_tokenService.TryDecryptToken(token, out SignalModel signal)) - return false; + return Enumerable.Empty(); - await TriggerSignalAsync(signal.Name, input, signal.WorkflowInstanceId, cancellationToken: cancellationToken); - return true; + return await TriggerSignalAsync(signal.Name, input, signal.WorkflowInstanceId, cancellationToken: cancellationToken); } - public async Task TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default) + public async Task> TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default) { - await _workflowLaunchpad.CollectAndExecuteWorkflowsAsync(new CollectWorkflowsContext( + return await _workflowLaunchpad.CollectAndExecuteWorkflowsAsync(new CollectWorkflowsContext( nameof(SignalReceived), new SignalReceivedBookmark { Signal = signal, WorkflowInstanceId = workflowInstanceId }, new SignalReceivedBookmark { Signal = signal }, @@ -43,16 +43,15 @@ namespace Elsa.Activities.Signaling.Services ), new Signal(signal, input), cancellationToken); } - public async Task DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default) + public async Task> DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default) { if (!_tokenService.TryDecryptToken(token, out SignalModel signal)) - return false; + return Enumerable.Empty(); - await DispatchSignalAsync(signal.Name, input, signal.WorkflowInstanceId, cancellationToken: cancellationToken); - return true; + return await DispatchSignalAsync(signal.Name, input, signal.WorkflowInstanceId, cancellationToken: cancellationToken); } - public async Task DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default) => + public async Task> DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default) => await _workflowLaunchpad.CollectAndDispatchWorkflowsAsync(new CollectWorkflowsContext( nameof(SignalReceived), new SignalReceivedBookmark { Signal = signal, WorkflowInstanceId = workflowInstanceId }, diff --git a/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs b/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs index 11f9f597f..e06a680c6 100644 --- a/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs +++ b/src/core/Elsa.Core/Services/WorkflowLaunchpad.cs @@ -219,11 +219,11 @@ namespace Elsa.Services return pendingWorkflow; } - public async Task> CollectAndExecuteWorkflowsAsync(CollectWorkflowsContext context, object? input = default, CancellationToken cancellationToken = default) + public async Task> CollectAndExecuteWorkflowsAsync(CollectWorkflowsContext context, object? input = default, CancellationToken cancellationToken = default) { var pendingWorkflows = await CollectWorkflowsAsync(context, cancellationToken).ToList(); await ExecutePendingWorkflowsAsync(pendingWorkflows, input, cancellationToken); - return pendingWorkflows; + return pendingWorkflows.Select(x => new StartedWorkflow(x.WorkflowInstanceId, x.ActivityId)); } public async Task> CollectAndDispatchWorkflowsAsync(CollectWorkflowsContext context, object? input = default, CancellationToken cancellationToken = default) diff --git a/src/server/Elsa.Server.Api/Endpoints/Signals/Dispatch.cs b/src/server/Elsa.Server.Api/Endpoints/Signals/Dispatch.cs new file mode 100644 index 000000000..8dea7337f --- /dev/null +++ b/src/server/Elsa.Server.Api/Endpoints/Signals/Dispatch.cs @@ -0,0 +1,47 @@ +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.Signaling.Services; +using Elsa.Server.Api.ActionFilters; +using Elsa.Services.Models; +using Microsoft.AspNetCore.Http; +using Microsoft.AspNetCore.Mvc; +using Open.Linq.AsyncExtensions; +using Swashbuckle.AspNetCore.Annotations; + +namespace Elsa.Server.Api.Endpoints.Signals +{ + [ApiController] + [ApiVersion("1")] + [Route("v{apiVersion:apiVersion}/signals/{signalName}/dispatch")] + [Produces("application/json")] + public class Dispatch : Controller + { + private readonly ISignaler _signaler; + + public Dispatch(ISignaler signaler) + { + _signaler = signaler; + } + + [HttpPost] + [ElsaJsonFormatter] + [ProducesResponseType(StatusCodes.Status200OK, Type = typeof(DispatchSignalResponse))] + [SwaggerOperation( + Summary = "Signals all workflows waiting on the specified signal name synchronously.", + Description = "Signals all workflows waiting on the specified signal name synchronously.", + OperationId = "Signals.Execute", + Tags = new[] { "Signals" }) + ] + public async Task Handle(string signalName, DispatchSignalRequest request, CancellationToken cancellationToken = default) + { + var result = await _signaler.DispatchSignalAsync(signalName, request.Input, request.WorkflowInstanceId, request.CorrelationId, cancellationToken).ToList(); + + if (Response.HasStarted) + return new EmptyResult(); + + return Ok(new DispatchSignalResponse(result)); + } + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Api/Endpoints/Signals/Execute.cs b/src/server/Elsa.Server.Api/Endpoints/Signals/Execute.cs new file mode 100644 index 000000000..372af9005 --- /dev/null +++ b/src/server/Elsa.Server.Api/Endpoints/Signals/Execute.cs @@ -0,0 +1,47 @@ +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.Signaling.Services; +using Elsa.Server.Api.ActionFilters; +using Elsa.Services.Models; +using Microsoft.AspNetCore.Http; +using Microsoft.AspNetCore.Mvc; +using Open.Linq.AsyncExtensions; +using Swashbuckle.AspNetCore.Annotations; + +namespace Elsa.Server.Api.Endpoints.Signals +{ + [ApiController] + [ApiVersion("1")] + [Route("v{apiVersion:apiVersion}/signals/{signalName}/execute")] + [Produces("application/json")] + public class Execute : Controller + { + private readonly ISignaler _signaler; + + public Execute(ISignaler signaler) + { + _signaler = signaler; + } + + [HttpPost] + [ElsaJsonFormatter] + [ProducesResponseType(StatusCodes.Status200OK, Type = typeof(ExecuteSignalResponse))] + [SwaggerOperation( + Summary = "Signals all workflows waiting on the specified signal name synchronously.", + Description = "Signals all workflows waiting on the specified signal name synchronously.", + OperationId = "Signals.Execute", + Tags = new[] { "Signals" }) + ] + public async Task Handle(string signalName, ExecuteSignalRequest request, CancellationToken cancellationToken = default) + { + var result = await _signaler.TriggerSignalAsync(signalName, request.Input, request.WorkflowInstanceId, request.CorrelationId, cancellationToken).ToList(); + + if (Response.HasStarted) + return new EmptyResult(); + + return Ok(new ExecuteSignalResponse(result.Select(x => new StartedWorkflow(x.WorkflowInstanceId, x.ActivityId)).ToList())); + } + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Api/Endpoints/Signals/Models.cs b/src/server/Elsa.Server.Api/Endpoints/Signals/Models.cs new file mode 100644 index 000000000..ef191901e --- /dev/null +++ b/src/server/Elsa.Server.Api/Endpoints/Signals/Models.cs @@ -0,0 +1,14 @@ +using System.Collections.Generic; +using Elsa.Services.Models; +using Microsoft.AspNetCore.Mvc; + +namespace Elsa.Server.Api.Endpoints.Signals +{ + public record DispatchSignalRequest(string? WorkflowInstanceId, string? CorrelationId, object? Input); + + public record DispatchSignalResponse(ICollection StartedWorkflows); + + public record ExecuteSignalRequest(string? WorkflowInstanceId, string? CorrelationId, object? Input); + + public record ExecuteSignalResponse(ICollection StartedWorkflows); +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Api/Endpoints/WorkflowDefinitions/Models.cs b/src/server/Elsa.Server.Api/Endpoints/WorkflowDefinitions/Models.cs index e96a22b1e..7dad56196 100644 --- a/src/server/Elsa.Server.Api/Endpoints/WorkflowDefinitions/Models.cs +++ b/src/server/Elsa.Server.Api/Endpoints/WorkflowDefinitions/Models.cs @@ -35,7 +35,7 @@ namespace Elsa.Server.Api.Endpoints.WorkflowDefinitions public record ExecuteWorkflowsRequest(string ActivityType, IBookmark? Bookmark, IBookmark? Trigger, string? CorrelationId, string? WorkflowInstanceId, string? ContextId, object? Input); - public record ExecuteWorkflowsResponse(ICollection PendingWorkflows); + public record ExecuteWorkflowsResponse(ICollection StartedWorkflows); public record DispatchWorkflowsRequest(string ActivityType, IBookmark? Bookmark, IBookmark? Trigger, string? CorrelationId, string? WorkflowInstanceId, string? ContextId, object? Input); diff --git a/src/server/Elsa.Server.Api/Endpoints/Workflows/Dispatch.cs b/src/server/Elsa.Server.Api/Endpoints/Workflows/Dispatch.cs index ad2b696d7..c61518528 100644 --- a/src/server/Elsa.Server.Api/Endpoints/Workflows/Dispatch.cs +++ b/src/server/Elsa.Server.Api/Endpoints/Workflows/Dispatch.cs @@ -41,7 +41,7 @@ namespace Elsa.Server.Api.Endpoints.Workflows if (Response.HasStarted) return new EmptyResult(); - return Ok(new ExecuteWorkflowsResponse(result)); + return Ok(new DispatchWorkflowsResponse(result)); } } } \ No newline at end of file