Fix signaler and TriggerWorkflows handler

This commit is contained in:
Sipke Schoorstra 2021-05-08 20:44:09 +02:00
parent ab8c092887
commit 5761ad7145
10 changed files with 67 additions and 89 deletions

View file

@ -29,8 +29,8 @@ namespace Elsa.Activities.Http.Endpoints.Signals
await _signaler.TriggerSignalAsync(signal.Name, null, signal.WorkflowInstanceId, cancellationToken);
return HttpContext.Items.ContainsKey(WorkflowHttpResult.Instance)
? (IActionResult)new EmptyResult()
return HttpContext.Response.HasStarted
? new EmptyResult()
: Accepted();
}
}

View file

@ -24,7 +24,6 @@ namespace Elsa.Activities.Http.Results
var httpContext = _httpContextAccessor.HttpContext;
var response = httpContext.Response;
httpContext.Items[WorkflowHttpResult.Instance] = WorkflowHttpResult.Instance;
response.Redirect(Location.ToString(), Permanent);
}
}

View file

@ -1,7 +0,0 @@
namespace Elsa.Activities.Http
{
public class WorkflowHttpResult
{
public static readonly WorkflowHttpResult Instance = new WorkflowHttpResult();
}
}

View file

@ -3,7 +3,7 @@ using MediatR;
namespace Elsa.Dispatch
{
public record TriggerWorkflowsRequest(string ActivityType, IBookmark Bookmark, IBookmark Trigger, object? Input = default, string? CorrelationId = default, string? ContextId = default, string? TenantId = default) : IRequest<int>;
public record TriggerWorkflowsRequest(string ActivityType, IBookmark Bookmark, IBookmark Trigger, object? Input = default, string? CorrelationId = default, string? WorkflowInstanceId = default, string? ContextId = default, string? TenantId = default) : IRequest<int>;
public record ExecuteWorkflowDefinitionRequest(string WorkflowDefinitionId, string? ActivityId = default, object? Input = default, string? CorrelationId = default, string? ContextId = default, string? TenantId = default) : IRequest<Unit>;
public record ExecuteWorkflowInstanceRequest(string WorkflowInstanceId, string ActivityId, object? Input = default) : IRequest<Unit>;

View file

@ -20,17 +20,20 @@ namespace Elsa.Activities.Signaling.Services
_workflowDispatcher = workflowDispatcher;
}
public async Task TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, CancellationToken cancellationToken = default) =>
public async Task TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, CancellationToken cancellationToken = default)
{
await _mediator.Send(new TriggerWorkflowsRequest(
nameof(SignalReceived),
new SignalReceivedBookmark {Signal = signal, WorkflowInstanceId = workflowInstanceId},
new SignalReceivedBookmark {Signal = signal},
new SignalReceivedBookmark { Signal = signal, WorkflowInstanceId = workflowInstanceId },
new SignalReceivedBookmark { Signal = signal },
new Signal(signal, input),
default,
workflowInstanceId,
default,
TenantId),
cancellationToken
);
}
public async Task DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, CancellationToken cancellationToken = default) =>
await _workflowDispatcher.DispatchAsync(new TriggerWorkflowsRequest(

View file

@ -22,8 +22,6 @@ namespace Elsa.Dispatch.Handlers
private readonly IBookmarkFinder _bookmarkFinder;
private readonly ITriggerFinder _triggerFinder;
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly IWorkflowDefinitionDispatcher _workflowDefinitionDispatcher;
private readonly IWorkflowInstanceDispatcher _workflowInstanceDispatcher;
private readonly IMediator _mediator;
private readonly ElsaOptions _elsaOptions;
private readonly ILogger<TriggerWorkflows> _logger;
@ -33,8 +31,6 @@ namespace Elsa.Dispatch.Handlers
IBookmarkFinder bookmarkFinder,
ITriggerFinder triggerFinder,
IDistributedLockProvider distributedLockProvider,
IWorkflowDefinitionDispatcher workflowDefinitionDispatcher,
IWorkflowInstanceDispatcher workflowInstanceDispatcher,
IMediator mediator,
ElsaOptions elsaOptions,
ILogger<TriggerWorkflows> logger)
@ -43,8 +39,6 @@ namespace Elsa.Dispatch.Handlers
_bookmarkFinder = bookmarkFinder;
_triggerFinder = triggerFinder;
_distributedLockProvider = distributedLockProvider;
_workflowDefinitionDispatcher = workflowDefinitionDispatcher;
_workflowInstanceDispatcher = workflowInstanceDispatcher;
_mediator = mediator;
_elsaOptions = elsaOptions;
_logger = logger;
@ -55,31 +49,62 @@ namespace Elsa.Dispatch.Handlers
var correlationId = request.CorrelationId;
if (!string.IsNullOrWhiteSpace(correlationId))
return await ResumeCorrelatedWorkflowsAsync(request, cancellationToken);
if (!string.IsNullOrWhiteSpace(request.WorkflowInstanceId))
return await ResumeSpecificWorkflowInstanceAsync(request, cancellationToken);
return await TriggerWorkflowsAsync(request, cancellationToken);
}
private async Task<int> TriggerWorkflowsAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken)
{
var bookmarkResultsQuery = await _bookmarkFinder.FindBookmarksAsync(request.ActivityType, request.Bookmark, request.TenantId, cancellationToken);
var bookmarkResults = bookmarkResultsQuery.ToList();
var triggeredCount = bookmarkResults.GroupBy(x => x.WorkflowInstanceId).Select(x => x.Key).Distinct().Count();
await ResumeWorkflowsAsync(bookmarkResults, request.Input, cancellationToken);
var startedCount = await StartWorkflowsAsync(request, cancellationToken);
return startedCount + triggeredCount;
}
private async Task<int> ResumeSpecificWorkflowInstanceAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken)
{
var bookmarkResultsQuery = await _bookmarkFinder.FindBookmarksAsync(request.ActivityType, request.Bookmark, request.TenantId, cancellationToken);
bookmarkResultsQuery = bookmarkResultsQuery.Where(x => x.WorkflowInstanceId == request.WorkflowInstanceId);
var bookmarkResults = bookmarkResultsQuery.ToList();
var triggeredCount = bookmarkResults.GroupBy(x => x.WorkflowInstanceId).Select(x => x.Key).Distinct().Count();
await ResumeWorkflowsAsync(bookmarkResults, request.Input, cancellationToken);
return triggeredCount;
}
private async Task<int> ResumeCorrelatedWorkflowsAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken)
{
var correlationId = request.CorrelationId!;
var lockKey = correlationId;
_logger.LogDebug("Acquiring lock on correlation ID {CorrelationId}", correlationId);
await using var handle = await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken);
if (handle == null)
throw new LockAcquisitionException($"Failed to acquire a lock on {lockKey}");
var correlatedWorkflowInstanceCount = !string.IsNullOrWhiteSpace(correlationId)
? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(correlationId).WithStatus(WorkflowStatus.Suspended), cancellationToken)
: 0;
_logger.LogDebug("Found {CorrelatedWorkflowCount} correlated workflows,", correlatedWorkflowInstanceCount);
if (correlatedWorkflowInstanceCount > 0)
{
var lockKey = correlationId;
_logger.LogDebug("Acquiring lock on correlation ID {CorrelationId}", correlationId);
await using var handle = await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken);
if (handle == null)
throw new LockAcquisitionException($"Failed to acquire a lock on {lockKey}");
var correlatedWorkflowInstanceCount = !string.IsNullOrWhiteSpace(correlationId)
? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(correlationId).WithStatus(WorkflowStatus.Suspended), cancellationToken)
: 0;
_logger.LogDebug("Found {CorrelatedWorkflowCount} correlated workflows,", correlatedWorkflowInstanceCount);
if (correlatedWorkflowInstanceCount > 0)
{
_logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId);
var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(request.ActivityType, request.Bookmark, request.TenantId, cancellationToken).ToList();
await ResumeWorkflowsAsync(bookmarkResults, request.Input, cancellationToken);
return correlatedWorkflowInstanceCount;
}
_logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId);
var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(request.ActivityType, request.Bookmark, request.TenantId, cancellationToken).ToList();
await ResumeWorkflowsAsync(bookmarkResults, request.Input, cancellationToken);
}
return await StartWorkflowsAsync(request, cancellationToken);
return correlatedWorkflowInstanceCount;
}
private async Task<int> StartWorkflowsAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken)

View file

@ -14,8 +14,8 @@
</head>
<body class="h-screen" style="background-size: 30px 30px; background-image: url(/build/assets/images/tile.png); background-color: #FBFBFB;">
<elsa-studio-root server-url="https://localhost:11000/" monaco-lib-path="build/assets/js"></elsa-studio-root>
<!--<elsa-studio-root server-url="https://localhost:6080/" monaco-lib-path="build/assets/js"></elsa-studio-root>-->
<elsa-studio-root server-url="https://localhost:11000" monaco-lib-path="build/assets/js"></elsa-studio-root>
<!--<elsa-studio-root server-url="https://localhost:6080" monaco-lib-path="build/assets/js"></elsa-studio-root>-->
<!--<elsa-studio-root server-url="https://skynet-workflow.azurewebsites.net" monaco-lib-path="build/assets/js"></elsa-studio-root>-->
</body>

View file

@ -8,10 +8,6 @@
<AutoGenerateBindingRedirects>true</AutoGenerateBindingRedirects>
</PropertyGroup>
<ItemGroup>
<None Remove="elsa.db" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="Microsoft.VisualStudio.Azure.Containers.Tools.Targets" Version="1.10.9" />
</ItemGroup>

View file

@ -2,15 +2,11 @@ using Elsa.Activities.UserTask.Extensions;
using Elsa.Persistence.EntityFramework.Core.Extensions;
using Elsa.Persistence.EntityFramework.Sqlite;
using Elsa.Samples.Server.Host.Activities;
using Elsa.Server.Hangfire.Extensions;
using Hangfire;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Hosting;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NodaTime;
using NodaTime.Serialization.JsonNet;
namespace Elsa.Samples.Server.Host
{
@ -28,51 +24,22 @@ namespace Elsa.Samples.Server.Host
public void ConfigureServices(IServiceCollection services)
{
var elsaSection = Configuration.GetSection("Elsa");
var hangfireSection = Configuration.GetSection("Hangfire");
services.AddControllers();
// Hangfire is required when using Hangfire Dispatchers.
services
.AddHangfire(configuration => configuration
.UseSimpleAssemblyNameTypeSerializer()
.UseRecommendedSerializerSettings(settings => settings.ConfigureForNodaTime(DateTimeZoneProviders.Tzdb))
.UseInMemoryStorage())
.AddHangfireServer(options =>
{
options.ConfigureForElsaDispatchers();
hangfireSection
.GetSection("Server")
.Bind(options);
});
services
.AddActivityPropertyOptionsProvider<VehicleActivity>()
.AddRuntimeSelectItemsProvider<VehicleActivity>()
.AddElsa(elsa => elsa
.UseEntityFrameworkPersistence(ef => ef.UseSqlite())
//.UseEntityFrameworkPersistence(ef => ef.UseSqlServer("Server=LAPTOP-B76STK67;Database=Elsa;Integrated Security=true;MultipleActiveResultSets=True;"))
//.UseYesSqlPersistence()
// Using Hangfire as the dispatcher for workflow execution in the background.
.UseHangfireDispatchers()
.AddConsoleActivities()
.AddHttpActivities(elsaSection.GetSection("Http").Bind)
.AddEmailActivities(elsaSection.GetSection("Smtp").Bind)
.AddQuartzTemporalActivities()
// .AddQuartzTemporalActivities(configureQuartz: quartz => quartz.UsePersistentStore(x =>
// {
// x.UseJsonSerializer();
// x.UseGenericDatabase("SqlServer", ado =>
// {
// ado.ConnectionString = "Server=LAPTOP-B76STK67;Database=Elsa;Integrated Security=true;MultipleActiveResultSets=True;";
// });
// }))
.AddJavaScriptActivities()
.AddUserTaskActivities()
.AddActivitiesFrom<VehicleActivity>()
.AddActivitiesFrom<Startup>()
.AddWorkflowsFrom<Startup>()
);
// Elsa API endpoints.

View file

@ -10,17 +10,12 @@
"AllowedHosts": "*",
"Elsa": {
"Http": {
"BaseUrl": "http://localhost:16398"
"BaseUrl": "https://localhost:11000"
},
"Smtp": {
"Host": "localhost",
"Port": "2525",
"DefaultSender": "noreply@acme.com"
}
},
"Hangfire": {
"Server": {
"WorkerCount": 10
}
}
}