From 5761ad7145ed2525efe7d10577a8b2d0bf4eeeb2 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 8 May 2021 20:44:09 +0200 Subject: [PATCH] Fix signaler and TriggerWorkflows handler --- .../Endpoints/Signals/TriggerEndpoint.cs | 4 +- .../Results/RedirectResult.cs | 1 - .../WorkflowHttpResult.cs | 7 -- src/core/Elsa.Abstractions/Dispatch/Models.cs | 2 +- .../Activities/Signaling/Services/Signaler.cs | 9 ++- .../Dispatch/Handlers/TriggerWorkflows.cs | 81 ++++++++++++------- .../elsa-workflows-studio/src/index.html | 4 +- ...shboard.Samples.AspNetCore.Monolith.csproj | 4 - .../Elsa.Samples.Server.Host/Startup.cs | 37 +-------- .../Elsa.Samples.Server.Host/appsettings.json | 7 +- 10 files changed, 67 insertions(+), 89 deletions(-) delete mode 100644 src/activities/Elsa.Activities.Http/WorkflowHttpResult.cs diff --git a/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs b/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs index 8a683ea6d..869bb4e45 100644 --- a/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs +++ b/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs @@ -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(); } } diff --git a/src/activities/Elsa.Activities.Http/Results/RedirectResult.cs b/src/activities/Elsa.Activities.Http/Results/RedirectResult.cs index 024441b91..1dcea6ea0 100644 --- a/src/activities/Elsa.Activities.Http/Results/RedirectResult.cs +++ b/src/activities/Elsa.Activities.Http/Results/RedirectResult.cs @@ -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); } } diff --git a/src/activities/Elsa.Activities.Http/WorkflowHttpResult.cs b/src/activities/Elsa.Activities.Http/WorkflowHttpResult.cs deleted file mode 100644 index 55b8ecd59..000000000 --- a/src/activities/Elsa.Activities.Http/WorkflowHttpResult.cs +++ /dev/null @@ -1,7 +0,0 @@ -namespace Elsa.Activities.Http -{ - public class WorkflowHttpResult - { - public static readonly WorkflowHttpResult Instance = new WorkflowHttpResult(); - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Dispatch/Models.cs b/src/core/Elsa.Abstractions/Dispatch/Models.cs index dd977601c..1647748c4 100644 --- a/src/core/Elsa.Abstractions/Dispatch/Models.cs +++ b/src/core/Elsa.Abstractions/Dispatch/Models.cs @@ -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; + 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; public record ExecuteWorkflowDefinitionRequest(string WorkflowDefinitionId, string? ActivityId = default, object? Input = default, string? CorrelationId = default, string? ContextId = default, string? TenantId = default) : IRequest; public record ExecuteWorkflowInstanceRequest(string WorkflowInstanceId, string ActivityId, object? Input = default) : IRequest; diff --git a/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs b/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs index 52669b27c..419a5cb45 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs @@ -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( diff --git a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs index 1968a0efe..7ecb8917a 100644 --- a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs +++ b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs @@ -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 _logger; @@ -33,8 +31,6 @@ namespace Elsa.Dispatch.Handlers IBookmarkFinder bookmarkFinder, ITriggerFinder triggerFinder, IDistributedLockProvider distributedLockProvider, - IWorkflowDefinitionDispatcher workflowDefinitionDispatcher, - IWorkflowInstanceDispatcher workflowInstanceDispatcher, IMediator mediator, ElsaOptions elsaOptions, ILogger 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 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 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 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(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(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 StartWorkflowsAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken) diff --git a/src/designer/elsa-workflows-studio/src/index.html b/src/designer/elsa-workflows-studio/src/index.html index b6402d31f..600e4680f 100644 --- a/src/designer/elsa-workflows-studio/src/index.html +++ b/src/designer/elsa-workflows-studio/src/index.html @@ -14,8 +14,8 @@ - - + + diff --git a/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj b/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj index 54b765d36..8f161359a 100644 --- a/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj +++ b/src/samples/dashboard/aspnetcore/ElsaDashboard.Samples.AspNetCore.Monolith/ElsaDashboard.Samples.AspNetCore.Monolith.csproj @@ -8,10 +8,6 @@ true - - - - diff --git a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs index 0c322a200..c90ab8dcf 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs +++ b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs @@ -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() .AddRuntimeSelectItemsProvider() .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() + .AddActivitiesFrom() + .AddWorkflowsFrom() ); // Elsa API endpoints. diff --git a/src/samples/server/Elsa.Samples.Server.Host/appsettings.json b/src/samples/server/Elsa.Samples.Server.Host/appsettings.json index 02e8092f2..597dddbc4 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/appsettings.json +++ b/src/samples/server/Elsa.Samples.Server.Host/appsettings.json @@ -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 - } } }