diff --git a/src/modules/Elsa.Alterations/Services/BackgroundAlterationJobDispatcher.cs b/src/modules/Elsa.Alterations/Services/BackgroundAlterationJobDispatcher.cs index e9180a566..0b426dc5a 100644 --- a/src/modules/Elsa.Alterations/Services/BackgroundAlterationJobDispatcher.cs +++ b/src/modules/Elsa.Alterations/Services/BackgroundAlterationJobDispatcher.cs @@ -1,4 +1,5 @@ using Elsa.Alterations.Core.Contracts; +using Elsa.Common.Multitenancy; using Elsa.Mediator.Contracts; using Microsoft.Extensions.DependencyInjection; @@ -7,19 +8,23 @@ namespace Elsa.Alterations.Services; /// /// Dispatches an alteration job for execution using an in-memory channel. /// -public class BackgroundAlterationJobDispatcher(IJobQueue jobQueue, IServiceScopeFactory scopeFactory) : IAlterationJobDispatcher +public class BackgroundAlterationJobDispatcher( + IJobQueue jobQueue, + ITenantAccessor tenantAccessor, + ITenantScopeFactory tenantScopeFactory) : IAlterationJobDispatcher { /// public ValueTask DispatchAsync(string jobId, CancellationToken cancellationToken = default) { - jobQueue.Enqueue(ct => ExecuteJobAsync(jobId, ct)); + var tenant = tenantAccessor.Tenant; + jobQueue.Enqueue(ct => ExecuteJobAsync(jobId, tenant, ct)); return default; } - private async Task ExecuteJobAsync(string alterationJobId, CancellationToken cancellationToken) + private async Task ExecuteJobAsync(string alterationJobId, Tenant? tenant, CancellationToken cancellationToken) { - using var scope = scopeFactory.CreateScope(); - var alterationJobRunner = scope.ServiceProvider.GetRequiredService(); + await using var tenantScope = tenantScopeFactory.CreateScope(tenant); + var alterationJobRunner = tenantScope.ServiceProvider.GetRequiredService(); await alterationJobRunner.RunAsync(alterationJobId, cancellationToken); } } diff --git a/test/integration/Elsa.Alterations.IntegrationTests/BackgroundAlterationJobDispatcherTests.cs b/test/integration/Elsa.Alterations.IntegrationTests/BackgroundAlterationJobDispatcherTests.cs new file mode 100644 index 000000000..4c3fa73e0 --- /dev/null +++ b/test/integration/Elsa.Alterations.IntegrationTests/BackgroundAlterationJobDispatcherTests.cs @@ -0,0 +1,164 @@ +using System.Collections.Concurrent; +using Elsa.Alterations.Core.Contracts; +using Elsa.Alterations.Core.Entities; +using Elsa.Alterations.Services; +using Elsa.Common.Multitenancy; +using Elsa.Mediator.Contracts; +using Microsoft.Extensions.DependencyInjection; +using NSubstitute; + +namespace Elsa.Alterations.IntegrationTests; + +public class BackgroundAlterationJobDispatcherTests : IAsyncLifetime +{ + private readonly List> _queuedCallbacks = []; + private readonly DefaultTenantAccessor _tenantAccessor; + private readonly RecordingAlterationJobRunner _runner; + private readonly ServiceProvider _serviceProvider; + + public BackgroundAlterationJobDispatcherTests() + { + var jobQueue = Substitute.For(); + jobQueue + .Enqueue(Arg.Any>()) + .Returns(callInfo => + { + var callback = callInfo.Arg>(); + _queuedCallbacks.Add(callback); + return $"queued-job-{_queuedCallbacks.Count}"; + }); + + _tenantAccessor = new DefaultTenantAccessor(); + _runner = new RecordingAlterationJobRunner(_tenantAccessor); + var services = new ServiceCollection() + .AddSingleton(jobQueue) + .AddSingleton(_tenantAccessor) + .AddSingleton() + .AddScoped(_ => _runner) + .AddScoped(); + _serviceProvider = services.BuildServiceProvider(validateScopes: true); + } + + public Task InitializeAsync() => Task.CompletedTask; + + public async Task DisposeAsync() => await _serviceProvider.DisposeAsync(); + + [Fact] + public async Task DispatchAsync_WhenQueuedWorkRunsAfterDispatchScopeEnds_PreservesTenant() + { + const string jobId = "alteration-job"; + var dispatchingTenant = new Tenant { Id = "tenant-a", Name = "Tenant A" }; + var workerTenant = new Tenant { Id = "tenant-b", Name = "Tenant B" }; + + await DispatchAndRunUnderWorkerTenantAsync(jobId, dispatchingTenant, workerTenant, callback => callback(CancellationToken.None)); + } + + [Fact] + public async Task DispatchAsync_WhenRunnerThrows_StillRestoresWorkerTenant() + { + const string jobId = "alteration-job"; + var dispatchingTenant = new Tenant { Id = "tenant-a", Name = "Tenant A" }; + var workerTenant = new Tenant { Id = "tenant-b", Name = "Tenant B" }; + _runner.ExceptionToThrow = new InvalidOperationException("Runner failure"); + + await DispatchAndRunUnderWorkerTenantAsync( + jobId, + dispatchingTenant, + workerTenant, + callback => Assert.ThrowsAsync(() => callback(CancellationToken.None))); + } + + [Fact] + public async Task DispatchAsync_WhenNoTenantIsPushedAtDispatchTime_UsesDefaultTenant() + { + const string jobId = "alteration-job"; + + await DispatchAsync(jobId, tenant: null); + + var callback = Assert.Single(_queuedCallbacks); + await callback(CancellationToken.None); + + // With no tenant pushed at dispatch time, DefaultTenantScopeFactory.CreateScope(null) pushes a null + // tenant onto the accessor. DefaultTenantAccessor.TenantId then falls back to Tenant.DefaultTenantId + // (an empty string) rather than null, so that is the value the runner observes. + Assert.Equal(Tenant.DefaultTenantId, _runner.GetObservedTenantId(jobId)); + Assert.Null(_tenantAccessor.Tenant); + } + + /// + /// Dispatches a job under , then runs the single queued callback while a + /// different is active on the accessor, invoking it via + /// so callers can assert success or failure. Asserts the tenant behavior + /// common to both outcomes: the worker tenant remains active for the duration of the callback, the runner + /// observed the dispatching tenant, and the accessor's tenant is restored to null afterward. + /// + private async Task DispatchAndRunUnderWorkerTenantAsync( + string jobId, + Tenant dispatchingTenant, + Tenant workerTenant, + Func, Task> runCallbackAsync) + { + await DispatchAsync(jobId, dispatchingTenant); + + var callback = Assert.Single(_queuedCallbacks); + using (_tenantAccessor.PushContext(workerTenant)) + { + await runCallbackAsync(callback); + Assert.Same(workerTenant, _tenantAccessor.Tenant); + } + + Assert.Equal(dispatchingTenant.Id, _runner.GetObservedTenantId(jobId)); + Assert.Null(_tenantAccessor.Tenant); + } + + [Fact] + public async Task DispatchAsync_WhenConcurrentJobsBelongToDifferentTenants_TenantsAreNotExchanged() + { + const string jobAId = "alteration-job-a"; + const string jobBId = "alteration-job-b"; + var tenantA = new Tenant { Id = "tenant-a", Name = "Tenant A" }; + var tenantB = new Tenant { Id = "tenant-b", Name = "Tenant B" }; + + await DispatchAsync(jobAId, tenantA); + await DispatchAsync(jobBId, tenantB); + + Assert.Equal(2, _queuedCallbacks.Count); + var callbackA = _queuedCallbacks[0]; + var callbackB = _queuedCallbacks[1]; + + await Task.WhenAll(callbackA(CancellationToken.None), callbackB(CancellationToken.None)); + + Assert.Equal(tenantA.Id, _runner.GetObservedTenantId(jobAId)); + Assert.Equal(tenantB.Id, _runner.GetObservedTenantId(jobBId)); + } + + private async Task DispatchAsync(string jobId, Tenant? tenant) + { + using (_tenantAccessor.PushContext(tenant)) + using (var dispatchScope = _serviceProvider.CreateScope()) + { + var dispatcher = dispatchScope.ServiceProvider.GetRequiredService(); + await dispatcher.DispatchAsync(jobId); + } + } + + private sealed class RecordingAlterationJobRunner(ITenantAccessor tenantAccessor) : IAlterationJobRunner + { + private readonly ConcurrentDictionary _observedTenantIdsByJobId = new(); + + public Exception? ExceptionToThrow { get; set; } + + public string? GetObservedTenantId(string jobId) => _observedTenantIdsByJobId.TryGetValue(jobId, out var tenantId) ? tenantId : null; + + public async Task RunAsync(string jobId, CancellationToken cancellationToken = default) + { + await Task.Yield(); + _observedTenantIdsByJobId[jobId] = tenantAccessor.TenantId; + + if (ExceptionToThrow is not null) + throw ExceptionToThrow; + + return new AlterationJob { Id = jobId }; + } + } +}