Refactor signals and activity cancellation logic

Changed static signals to instance variables in test classes for better encapsulation. Increased default timeout in `SignalManager` to 60 seconds. Removed unused cancellation token logic in `ActivityExecutionContext`.
This commit is contained in:
Sipke Schoorstra 2024-10-30 19:04:48 +01:00
parent 6be9e53de1
commit c49a92b5b7
6 changed files with 17 additions and 39 deletions

View file

@ -6,7 +6,7 @@ public class SignalManager
{
private readonly ConcurrentDictionary<object, TaskCompletionSource<object?>> _signals = new();
public async Task<T> WaitAsync<T>(object signal, int millisecondsTimeout = 10000)
public async Task<T> WaitAsync<T>(object signal, int millisecondsTimeout = 60000)
{
var result = await WaitAsync(signal, millisecondsTimeout);
@ -16,7 +16,7 @@ public class SignalManager
return typedResult;
}
public async Task<object?> WaitAsync(object signal, int millisecondsTimeout = 10000)
public async Task<object?> WaitAsync(object signal, int millisecondsTimeout = 60000)
{
var taskCompletionSource = GetOrCreate(signal);
using var cancellationTokenSource = new CancellationTokenSource(millisecondsTimeout);

View file

@ -7,19 +7,8 @@ namespace Elsa.Workflows;
public partial class ActivityExecutionContext
{
private readonly CancellationTokenRegistration _cancellationRegistration;
private readonly CancellationTokenSource _cancellationTokenSource;
private readonly INotificationSender _publisher;
private void CancelActivity()
{
// If the activity is not running, do nothing.
if (Status != ActivityStatus.Running && Status != ActivityStatus.Faulted)
return;
_ = Task.Run(async () => await CancelActivityAsync());
}
private bool CanCancelActivity()
{
return Status is not ActivityStatus.Canceled and not ActivityStatus.Completed;
@ -30,12 +19,6 @@ public partial class ActivityExecutionContext
if(!CanCancelActivity())
return;
// Select all child contexts.
var childContexts = WorkflowExecutionContext.ActivityExecutionContexts.Where(x => x.ParentActivityExecutionContext == this).ToList();
foreach (var childContext in childContexts)
childContext._cancellationTokenSource.Cancel();
TransitionTo(ActivityStatus.Canceled);
ClearBookmarks();
ClearCompletionCallbacks();
@ -44,7 +27,6 @@ public partial class ActivityExecutionContext
// Add an execution log entry.
AddExecutionLogEntry("Canceled", payload: JournalData);
await _cancellationRegistration.DisposeAsync();
await this.SendSignalAsync(new CancelSignal());
// ReSharper disable once MethodSupportsCancellation

View file

@ -49,9 +49,6 @@ public partial class ActivityExecutionContext : IExecutionContext, IDisposable
CancellationToken = cancellationToken;
Id = id;
_publisher = GetRequiredService<INotificationSender>();
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
_cancellationRegistration = _cancellationTokenSource.Token.Register(CancelActivity);
}
/// <summary>
@ -136,9 +133,6 @@ public partial class ActivityExecutionContext : IExecutionContext, IDisposable
public void TransitionTo(ActivityStatus status)
{
Status = status;
if (Status is ActivityStatus.Completed or ActivityStatus.Canceled or ActivityStatus.Faulted)
_cancellationRegistration.Dispose();
}
/// <summary>
@ -678,7 +672,5 @@ public partial class ActivityExecutionContext : IExecutionContext, IDisposable
void IDisposable.Dispose()
{
_cancellationRegistration.Dispose();
_cancellationTokenSource.Dispose();
}
}

View file

@ -15,6 +15,7 @@ using Elsa.MassTransit.Extensions;
using Elsa.Testing.Shared.Handlers;
using Elsa.Testing.Shared.Services;
using Elsa.Workflows.Management;
using Elsa.Workflows.Runtime.Distributed.Extensions;
using FluentStorage;
using Hangfire.Annotations;
using Microsoft.AspNetCore.Hosting;
@ -87,8 +88,11 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
runtime.UseEntityFrameworkCore(ef => ef.UsePostgreSql(dbConnectionString));
runtime.UseCache();
runtime.UseMassTransitDispatcher();
runtime.UseProtoActor();
//runtime.UseDistributedRuntime();
runtime.UseDistributedRuntime();
});
elsa.UseDistributedCache(distributedCaching =>
{
distributedCaching.UseMassTransit();
});
elsa.UseJavaScript(options =>
{

View file

@ -15,7 +15,7 @@ public class BulkDispatchWorkflowsTests : AppComponentTest
private readonly WorkflowEvents _workflowEvents;
private readonly SignalManager _signalManager;
private readonly IWorkflowRuntime _workflowRuntime;
private static readonly object GreetEmployeesWorkflowCompletedSignal = new();
private readonly object _greetEmployeesWorkflowCompletedSignal = new();
public BulkDispatchWorkflowsTests(App app) : base(app)
{
@ -35,7 +35,7 @@ public class BulkDispatchWorkflowsTests : AppComponentTest
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionId(GreetEmployeesWorkflow.DefinitionId, VersionOptions.Published)
});
await workflowClient.RunInstanceAsync(RunWorkflowInstanceRequest.Empty);
var parentWorkflowInstanceArgs = await _signalManager.WaitAsync<WorkflowInstanceSavedEventArgs>(GreetEmployeesWorkflowCompletedSignal);
var parentWorkflowInstanceArgs = await _signalManager.WaitAsync<WorkflowInstanceSavedEventArgs>(_greetEmployeesWorkflowCompletedSignal);
Assert.Equal(WorkflowStatus.Finished, parentWorkflowInstanceArgs.WorkflowInstance.Status);
}
@ -64,7 +64,7 @@ public class BulkDispatchWorkflowsTests : AppComponentTest
return;
if (e.WorkflowInstance.DefinitionId == GreetEmployeesWorkflow.DefinitionId)
_signalManager.Trigger(GreetEmployeesWorkflowCompletedSignal, e);
_signalManager.Trigger(_greetEmployeesWorkflowCompletedSignal, e);
}
protected override void OnDispose()

View file

@ -16,8 +16,8 @@ public class DispatchWorkflowsTests : AppComponentTest
private readonly SignalManager _signalManager;
private readonly IWorkflowRuntime _workflowRuntime;
private static readonly object ChildWorkflowCompletedSignal = new();
private static readonly object ParentWorkflowCompletedSignal = new();
private readonly object _childWorkflowCompletedSignal = new();
private readonly object _parentWorkflowCompletedSignal = new();
public DispatchWorkflowsTests(App app) : base(app)
{
@ -36,8 +36,8 @@ public class DispatchWorkflowsTests : AppComponentTest
WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionId(DispatchAndWaitWorkflow.DefinitionId, VersionOptions.Published)
});
await workflowClient.RunInstanceAsync(RunWorkflowInstanceRequest.Empty);
var childWorkflowInstanceArgs = await _signalManager.WaitAsync<WorkflowInstanceSavedEventArgs>(ChildWorkflowCompletedSignal);
var parentWorkflowInstanceArgs = await _signalManager.WaitAsync<WorkflowInstanceSavedEventArgs>(ParentWorkflowCompletedSignal);
var childWorkflowInstanceArgs = await _signalManager.WaitAsync<WorkflowInstanceSavedEventArgs>(_childWorkflowCompletedSignal);
var parentWorkflowInstanceArgs = await _signalManager.WaitAsync<WorkflowInstanceSavedEventArgs>(_parentWorkflowCompletedSignal);
Assert.Equal(WorkflowStatus.Finished, childWorkflowInstanceArgs.WorkflowInstance.Status);
Assert.Equal(WorkflowStatus.Finished, parentWorkflowInstanceArgs.WorkflowInstance.Status);
@ -49,10 +49,10 @@ public class DispatchWorkflowsTests : AppComponentTest
return;
if (e.WorkflowInstance.DefinitionId == ChildWorkflow.DefinitionId)
_signalManager.Trigger(ChildWorkflowCompletedSignal, e);
_signalManager.Trigger(_childWorkflowCompletedSignal, e);
if (e.WorkflowInstance.DefinitionId == DispatchAndWaitWorkflow.DefinitionId)
_signalManager.Trigger(ParentWorkflowCompletedSignal, e);
_signalManager.Trigger(_parentWorkflowCompletedSignal, e);
}
protected override void OnDispose()