From 0acc836f4bb49acbd59a9fa3d7faa7f8a588b8d6 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 2 Apr 2021 16:53:46 +0200 Subject: [PATCH] Fix signaling --- .../Endpoints/Signals/DispatchEndpoint.cs | 2 +- .../Endpoints/Signals/TriggerEndpoint.cs | 2 +- .../Elsa.Activities.Http/Models/Signal.cs | 6 +++--- .../SignalReceived/SignalReceivedBookmark.cs | 4 ++-- .../Activities/Signaling/Services/ISignaler.cs | 4 ++-- .../Activities/Signaling/Services/Signaler.cs | 14 ++++++-------- ...ExecutionContextForActivityBlueprintFactory.cs | 15 +++++++-------- .../Elsa.Samples.SignalingConsole/Program.cs | 2 +- 8 files changed, 23 insertions(+), 26 deletions(-) diff --git a/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs b/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs index e45b1babf..56d4328b1 100644 --- a/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs +++ b/src/activities/Elsa.Activities.Http/Endpoints/Signals/DispatchEndpoint.cs @@ -27,7 +27,7 @@ namespace Elsa.Activities.Http.Endpoints.Signals if (!_tokenService.TryDecryptToken(token, out Signal signal)) return NotFound(); - await _signaler.DispatchSignalAsync(signal.Name, null, signal.CorrelationId, cancellationToken); + await _signaler.DispatchSignalAsync(signal.Name, null, signal.WorkflowInstanceId, cancellationToken); return Accepted(); } } diff --git a/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs b/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs index b30a0cc99..8a683ea6d 100644 --- a/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs +++ b/src/activities/Elsa.Activities.Http/Endpoints/Signals/TriggerEndpoint.cs @@ -27,7 +27,7 @@ namespace Elsa.Activities.Http.Endpoints.Signals if (!_tokenService.TryDecryptToken(token, out Signal signal)) return NotFound(); - await _signaler.TriggerSignalAsync(signal.Name, null, signal.CorrelationId, cancellationToken); + await _signaler.TriggerSignalAsync(signal.Name, null, signal.WorkflowInstanceId, cancellationToken); return HttpContext.Items.ContainsKey(WorkflowHttpResult.Instance) ? (IActionResult)new EmptyResult() diff --git a/src/activities/Elsa.Activities.Http/Models/Signal.cs b/src/activities/Elsa.Activities.Http/Models/Signal.cs index 1c61b8f6a..b61b039a5 100644 --- a/src/activities/Elsa.Activities.Http/Models/Signal.cs +++ b/src/activities/Elsa.Activities.Http/Models/Signal.cs @@ -6,13 +6,13 @@ namespace Elsa.Activities.Http.Models { } - public Signal(string name, string correlationId) + public Signal(string name, string workflowInstanceId) { Name = name; - CorrelationId = correlationId; + WorkflowInstanceId = workflowInstanceId; } public string Name { get; set; } = default!; - public string CorrelationId { get; set; } = default!; + public string WorkflowInstanceId { get; set; } = default!; } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/Activities/SignalReceived/SignalReceivedBookmark.cs b/src/core/Elsa.Core/Activities/Signaling/Activities/SignalReceived/SignalReceivedBookmark.cs index d19dd786d..44f39e3e0 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Activities/SignalReceived/SignalReceivedBookmark.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Activities/SignalReceived/SignalReceivedBookmark.cs @@ -9,7 +9,7 @@ namespace Elsa.Activities.Signaling public class SignalReceivedBookmark : IBookmark { public string Signal { get; set; } = default!; - public string? CorrelationId { get; set; } + public string? WorkflowInstanceId { get; set; } } public class SignalReceivedBookmarkProvider : BookmarkProvider @@ -20,7 +20,7 @@ namespace Elsa.Activities.Signaling new SignalReceivedBookmark { Signal = (await context.Activity.GetPropertyValueAsync(x => x.Signal, cancellationToken))!, - CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId + WorkflowInstanceId = context.ActivityExecutionContext.WorkflowInstance.Id } }; } diff --git a/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs b/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs index de3c5509f..84bf84db3 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs @@ -8,11 +8,11 @@ namespace Elsa.Activities.Signaling.Services /// /// Runs all workflows that start with or are blocked on the activity. /// - Task TriggerSignalAsync(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default); + Task TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, CancellationToken cancellationToken = default); /// /// Dispatches all workflows that start with or are blocked on the activity. /// - Task DispatchSignalAsync(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default); + Task DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = 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 faf367246..861dec247 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs @@ -20,25 +20,23 @@ namespace Elsa.Activities.Signaling.Services _workflowDispatcher = workflowDispatcher; } - public async Task TriggerSignalAsync(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default) => + public async Task TriggerSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, CancellationToken cancellationToken = default) => await _triggersWorkflows.TriggerWorkflowsAsync( nameof(SignalReceived), - new SignalReceivedBookmark { Signal = signal, CorrelationId = correlationId }, + new SignalReceivedBookmark { Signal = signal, WorkflowInstanceId = workflowInstanceId }, new SignalReceivedBookmark { Signal = signal }, - correlationId, + default, new Signal(signal, input), tenantId: TenantId, cancellationToken: cancellationToken ); - public async Task DispatchSignalAsync(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default) => + public async Task DispatchSignalAsync(string signal, object? input = default, string? workflowInstanceId = default, CancellationToken cancellationToken = default) => await _workflowDispatcher.DispatchAsync(new TriggerWorkflowsRequest( nameof(SignalReceived), - new SignalReceivedBookmark { Signal = signal, CorrelationId = correlationId }, + new SignalReceivedBookmark { Signal = signal, WorkflowInstanceId = workflowInstanceId }, new SignalReceivedBookmark { Signal = signal }, - new Signal(signal, input), - correlationId, - TenantId: TenantId), + new Signal(signal, input)), cancellationToken); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/ActivityExecutionContextForActivityBlueprintFactory.cs b/src/core/Elsa.Core/Services/ActivityExecutionContextForActivityBlueprintFactory.cs index 02c80c835..be9926f81 100644 --- a/src/core/Elsa.Core/Services/ActivityExecutionContextForActivityBlueprintFactory.cs +++ b/src/core/Elsa.Core/Services/ActivityExecutionContextForActivityBlueprintFactory.cs @@ -9,11 +9,11 @@ namespace Elsa.Services /// public class ActivityExecutionContextForActivityBlueprintFactory : ICreatesActivityExecutionContextForActivityBlueprint { - readonly IServiceProvider serviceProvider; + private readonly IServiceProvider _serviceProvider; public ActivityExecutionContextForActivityBlueprintFactory(IServiceProvider serviceProvider) { - this.serviceProvider = serviceProvider ?? throw new ArgumentNullException(nameof(serviceProvider)); + _serviceProvider = serviceProvider ?? throw new ArgumentNullException(nameof(serviceProvider)); } /// @@ -23,11 +23,10 @@ namespace Elsa.Services /// A workflow execution context /// A cancellation token /// An activity execution context - public ActivityExecutionContext CreateActivityExecutionContext(IActivityBlueprint activityBlueprint, - WorkflowExecutionContext workflowExecutionContext, - CancellationToken cancellationToken) - { - return new ActivityExecutionContext(serviceProvider, workflowExecutionContext, activityBlueprint, null, false, cancellationToken); - } + public ActivityExecutionContext CreateActivityExecutionContext( + IActivityBlueprint activityBlueprint, + WorkflowExecutionContext workflowExecutionContext, + CancellationToken cancellationToken) => + new(_serviceProvider, workflowExecutionContext, activityBlueprint, null, false, cancellationToken); } } \ No newline at end of file diff --git a/src/samples/console/Elsa.Samples.SignalingConsole/Program.cs b/src/samples/console/Elsa.Samples.SignalingConsole/Program.cs index 5afa3fe70..aa68ba5c1 100644 --- a/src/samples/console/Elsa.Samples.SignalingConsole/Program.cs +++ b/src/samples/console/Elsa.Samples.SignalingConsole/Program.cs @@ -40,7 +40,7 @@ namespace Elsa.Samples.SignalingConsole // The workflows are now suspended at the red light. // Trigger a green light signal for the first car. var signaler = services.GetRequiredService(); - await signaler.TriggerSignalAsync("Green", correlationId: "Car 2"); + await signaler.TriggerSignalAsync("Green", workflowInstanceId: "Car 2"); // Notice that only the workflow correlated to the second car executed.