Implement signals dispatch & execute endpoints

This commit is contained in:
Sipke Schoorstra 2021-05-19 13:32:00 +02:00
parent e4f4375520
commit f80c0e1716
11 changed files with 139 additions and 31 deletions

View file

@ -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<IActionResult> 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);
}
}
}

View file

@ -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<IActionResult> 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);
}
}
}

View file

@ -87,7 +87,7 @@ namespace Elsa.Services
/// <summary>
/// Collects and executes workflows that are ready for execution. This takes into account both resumable (suspended) workflows as well as startable workflows.
/// </summary>
Task<IEnumerable<PendingWorkflow>> CollectAndExecuteWorkflowsAsync(CollectWorkflowsContext context, object? input = default, CancellationToken cancellationToken = default);
Task<IEnumerable<StartedWorkflow>> CollectAndExecuteWorkflowsAsync(CollectWorkflowsContext context, object? input = null, CancellationToken cancellationToken = default);
/// <summary>
/// Collects and dispatches workflows that are ready for execution. This takes into account both resumable (suspended) workflows as well as startable workflows.

View file

@ -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
/// <summary>
/// Runs all workflows that start with or are blocked on the <see cref="SignalReceived"/> activity.
/// </summary>
Task<bool> TriggerSignalTokenAsync(string signalToken, object? input = default, CancellationToken cancellationToken = default);
Task<IEnumerable<StartedWorkflow>> TriggerSignalTokenAsync(string signalToken, object? input = default, CancellationToken cancellationToken = default);
/// <summary>
/// Runs all workflows that start with or are blocked on the <see cref="SignalReceived"/> activity.
/// </summary>
Task TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<StartedWorkflow>> TriggerSignalAsync(string signal, object? input = null, string? workflowInstanceId = null, string? correlationId = null, CancellationToken cancellationToken = default);
/// <summary>
/// Dispatches all workflows that start with or are blocked on the <see cref="SignalReceived"/> activity.
/// </summary>
Task<bool> DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default);
Task<IEnumerable<PendingWorkflow>> DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default);
/// <summary>
/// Dispatches all workflows that start with or are blocked on the <see cref="SignalReceived"/> activity.
/// </summary>
Task DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default);
Task<IEnumerable<PendingWorkflow>> DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, string? correlationId = default, CancellationToken cancellationToken = default);
}
}

View file

@ -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<bool> TriggerSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default)
public async Task<IEnumerable<StartedWorkflow>> TriggerSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default)
{
if (!_tokenService.TryDecryptToken(token, out SignalModel signal))
return false;
return Enumerable.Empty<StartedWorkflow>();
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<IEnumerable<StartedWorkflow>> 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<bool> DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default)
public async Task<IEnumerable<PendingWorkflow>> DispatchSignalTokenAsync(string token, object? input = default, CancellationToken cancellationToken = default)
{
if (!_tokenService.TryDecryptToken(token, out SignalModel signal))
return false;
return Enumerable.Empty<PendingWorkflow>();
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<IEnumerable<PendingWorkflow>> 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 },

View file

@ -219,11 +219,11 @@ namespace Elsa.Services
return pendingWorkflow;
}
public async Task<IEnumerable<PendingWorkflow>> CollectAndExecuteWorkflowsAsync(CollectWorkflowsContext context, object? input = default, CancellationToken cancellationToken = default)
public async Task<IEnumerable<StartedWorkflow>> 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<IEnumerable<PendingWorkflow>> CollectAndDispatchWorkflowsAsync(CollectWorkflowsContext context, object? input = default, CancellationToken cancellationToken = default)

View file

@ -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<IActionResult> 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));
}
}
}

View file

@ -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<IActionResult> 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()));
}
}
}

View file

@ -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<PendingWorkflow> StartedWorkflows);
public record ExecuteSignalRequest(string? WorkflowInstanceId, string? CorrelationId, object? Input);
public record ExecuteSignalResponse(ICollection<StartedWorkflow> StartedWorkflows);
}

View file

@ -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<PendingWorkflow> PendingWorkflows);
public record ExecuteWorkflowsResponse(ICollection<StartedWorkflow> StartedWorkflows);
public record DispatchWorkflowsRequest(string ActivityType, IBookmark? Bookmark, IBookmark? Trigger, string? CorrelationId, string? WorkflowInstanceId, string? ContextId, object? Input);

View file

@ -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));
}
}
}