diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs index ef673127a..a9ba30ba1 100644 --- a/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs +++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs @@ -16,6 +16,7 @@ public class WorkflowInstanceSummary return new() { Id = workflowInstance.Id, + TenantId = workflowInstance.TenantId, DefinitionId = workflowInstance.DefinitionId, DefinitionVersionId = workflowInstance.DefinitionVersionId, Version = workflowInstance.Version, @@ -36,6 +37,7 @@ public class WorkflowInstanceSummary public static Expression> FromInstanceExpression() => workflowInstance => new() { Id = workflowInstance.Id, + TenantId = workflowInstance.TenantId, DefinitionId = workflowInstance.DefinitionId, DefinitionVersionId = workflowInstance.DefinitionVersionId, Version = workflowInstance.Version, @@ -53,6 +55,9 @@ public class WorkflowInstanceSummary /// The ID of the workflow instance. public string Id { get; set; } = null!; + /// The ID of the tenant that owns the workflow instance. + public string? TenantId { get; set; } + /// The ID of the workflow definition. public string DefinitionId { get; set; } = null!; diff --git a/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs b/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs index 5cb2cc2ba..288edb611 100644 --- a/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs +++ b/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs @@ -1,4 +1,5 @@ using Elsa.Common; +using Elsa.Common.Multitenancy; using Elsa.Common.RecurringTasks; using Elsa.Workflows.Management; using Elsa.Workflows.Management.Filters; @@ -12,11 +13,13 @@ namespace Elsa.Workflows.Runtime.Tasks; [SingleNodeTask] [UsedImplicitly] public class RestartInterruptedWorkflowsTask( - IWorkflowInstanceStore workflowInstanceStore, IWorkflowRestarter workflowRestarter, + IWorkflowInstanceStore workflowInstanceStore, + ILogger logger, IOptions options, ISystemClock systemClock, - ILogger logger) : RecurringTask + ITenantService? tenantService = null, + ITenantAccessor? tenantAccessor = null) : RecurringTask { /// public override async Task ExecuteAsync(CancellationToken cancellationToken) @@ -28,7 +31,26 @@ public class RestartInterruptedWorkflowsTask( logger.LogInformation("Restarting interrupted workflows."); await foreach (var workflowInstance in workflowInstances) { - await workflowRestarter.RestartWorkflowAsync(workflowInstance.Id, cancellationToken: cancellationToken); + try + { + var tenantId = workflowInstance.TenantId ?? string.Empty; + + if (tenantService != null && tenantAccessor != null && !string.IsNullOrWhiteSpace(tenantId) && tenantId != Tenant.AgnosticTenantId) + { + var tenant = await tenantService.FindAsync(tenantId, cancellationToken) ?? new Tenant { Id = tenantId, Name = tenantId }; + + using (tenantAccessor.PushContext(tenant)) + await workflowRestarter.RestartWorkflowAsync(workflowInstance.Id, cancellationToken: cancellationToken); + + continue; + } + + await workflowRestarter.RestartWorkflowAsync(workflowInstance.Id, cancellationToken: cancellationToken); + } + catch (Exception ex) + { + logger.LogError(ex, "Failed to restart interrupted workflow {WorkflowInstanceId}", workflowInstance.Id); + } } logger.LogInformation("Finished restarting interrupted workflows."); } diff --git a/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/RestartInterruptedWorkflowsTaskTests.cs b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/RestartInterruptedWorkflowsTaskTests.cs new file mode 100644 index 000000000..2b23f2c70 --- /dev/null +++ b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/RestartInterruptedWorkflowsTaskTests.cs @@ -0,0 +1,204 @@ +using Elsa.Common; +using Elsa.Common.Models; +using Elsa.Common.Multitenancy; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Runtime.Options; +using Elsa.Workflows.Runtime.Tasks; +using Microsoft.Extensions.Logging; +using NSubstitute; + +namespace Elsa.Workflows.Runtime.UnitTests.Services; + +public class RestartInterruptedWorkflowsTaskTests +{ + [Fact(DisplayName = "ExecuteAsync restarts each workflow within its tenant context")] + public async Task ExecuteAsync_RestartsWithinWorkflowTenantContext() + { + var workflowInstanceStore = Substitute.For(); + var tenantAccessor = new DefaultTenantAccessor(); + var tenantService = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", "tenant-a", now), + CreateWorkflowInstance("workflow-2", "tenant-b", now) + }; + var observedTenantIds = new List(); + + clock.UtcNow.Returns(now); + tenantService.FindAsync("tenant-a", Arg.Any()).Returns(new Tenant { Id = "tenant-a", Name = "Tenant A" }); + tenantService.FindAsync("tenant-b", Arg.Any()).Returns(new Tenant { Id = "tenant-b", Name = "Tenant B" }); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + workflowRestarter + .When(x => x.RestartWorkflowAsync(Arg.Any(), Arg.Any())) + .Do(_ => observedTenantIds.Add(tenantAccessor.Tenant?.Id)); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock, tenantService, tenantAccessor); + + await task.ExecuteAsync(CancellationToken.None); + + Assert.Equal(new[] { "tenant-a", "tenant-b" }, observedTenantIds); + Assert.Null(tenantAccessor.Tenant); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-1", Arg.Any()); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-2", Arg.Any()); + } + + [Theory(DisplayName = "ExecuteAsync does not push tenant context for default or agnostic tenant instances")] + [InlineData(null)] + [InlineData("")] + [InlineData("*")] + public async Task ExecuteAsync_DefaultOrAgnosticTenant_DoesNotPushContext(string? tenantId) + { + var workflowInstanceStore = Substitute.For(); + var tenantAccessor = new DefaultTenantAccessor(); + var tenantService = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", tenantId, now) + }; + var observedTenantIds = new List(); + + clock.UtcNow.Returns(now); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + workflowRestarter + .When(x => x.RestartWorkflowAsync(Arg.Any(), Arg.Any())) + .Do(_ => observedTenantIds.Add(tenantAccessor.Tenant?.Id)); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock, tenantService, tenantAccessor); + + await task.ExecuteAsync(CancellationToken.None); + + Assert.Equal(new string?[] { null }, observedTenantIds); + Assert.Null(tenantAccessor.Tenant); + await tenantService.DidNotReceive().FindAsync(Arg.Any(), Arg.Any()); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-1", Arg.Any()); + } + + [Fact(DisplayName = "ExecuteAsync continues restarting later workflows after a failure")] + public async Task ExecuteAsync_ContinuesAfterFailure() + { + var workflowInstanceStore = Substitute.For(); + var tenantAccessor = new DefaultTenantAccessor(); + var tenantService = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", "tenant-a", now), + CreateWorkflowInstance("workflow-2", "tenant-b", now) + }; + + clock.UtcNow.Returns(now); + tenantService.FindAsync("tenant-a", Arg.Any()).Returns(_ => Task.FromException(new InvalidOperationException("Transient tenant lookup failure"))); + tenantService.FindAsync("tenant-b", Arg.Any()).Returns(new Tenant { Id = "tenant-b", Name = "Tenant B" }); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock, tenantService, tenantAccessor); + + await task.ExecuteAsync(CancellationToken.None); + + await workflowRestarter.DidNotReceive().RestartWorkflowAsync("workflow-1", Arg.Any()); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-2", Arg.Any()); + } + + [Fact(DisplayName = "ExecuteAsync restarts tenant-specific workflows when tenant services are unavailable")] + public async Task ExecuteAsync_WithoutTenantServices_RestartsWorkflow() + { + var workflowInstanceStore = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", "tenant-a", now) + }; + + clock.UtcNow.Returns(now); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock); + + await task.ExecuteAsync(CancellationToken.None); + + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-1", Arg.Any()); + } + + private static WorkflowInstanceSummary CreateWorkflowInstance(string id, string? tenantId, DateTimeOffset updatedAt) + { + return new WorkflowInstanceSummary + { + Id = id, + TenantId = tenantId, + DefinitionId = "definition", + DefinitionVersionId = "definition:1", + Status = WorkflowStatus.Running, + SubStatus = WorkflowSubStatus.Pending, + CreatedAt = updatedAt.AddMinutes(-10), + UpdatedAt = updatedAt.AddMinutes(-10) + }; + } +} \ No newline at end of file