< Summary

Information
Class: Elsa.Workflows.Runtime.BookmarkQueueWorker
Assembly: Elsa.Workflows.Runtime
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/BookmarkQueueWorker.cs
Line coverage
93%
Covered lines: 56
Uncovered lines: 4
Coverable lines: 60
Total lines: 122
Line coverage: 93.3%
Branch coverage
92%
Covered branches: 26
Total branches: 28
Branch coverage: 92.8%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
get_Signaler()100%11100%
get_CapturedTenantId()100%22100%
.ctor(...)100%11100%
.ctor(...)75%44100%
Start()83.33%6685.71%
Stop()100%44100%
AwaitSignalAsync()100%4472.72%
InvokeProcessAsync(...)50%22100%
ProcessAsync()100%66100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/BookmarkQueueWorker.cs

#LineLine coverage
 1using Elsa.Common.Multitenancy;
 2using Microsoft.Extensions.DependencyInjection;
 3using Microsoft.Extensions.Logging;
 4using ThrottleDebounce;
 5
 6namespace Elsa.Workflows.Runtime;
 7
 8public class BookmarkQueueWorker : IBookmarkQueueWorker
 9{
 10    private readonly RateLimitedFunc<CancellationToken, Task>? _rateLimitedProcessAsync;
 11    private CancellationTokenSource _cts = null!;
 12    private bool _running;
 13    private readonly IServiceScopeFactory _scopeFactory;
 14    private readonly ITenantScopeFactory? _tenantScopeFactory;
 15    private readonly ITenantAccessor? _tenantAccessor;
 16    private readonly ILogger<BookmarkQueueWorker> _logger;
 17    private Tenant? _tenant;
 51918    protected IBookmarkQueueSignaler Signaler { get; }
 4119    protected internal string CapturedTenantId => (_tenant?.Id).NormalizeTenantId();
 20
 21    public BookmarkQueueWorker(
 22        IBookmarkQueueSignaler signaler,
 23        IServiceScopeFactory scopeFactory,
 24        ILogger<BookmarkQueueWorker> logger,
 25        ITenantScopeFactory? tenantScopeFactory = null,
 26        ITenantAccessor? tenantAccessor = null)
 8127        : this(signaler, scopeFactory, logger, tenantScopeFactory, tenantAccessor, TimeSpan.FromMilliseconds(500))
 28    {
 8129    }
 30
 11031    protected BookmarkQueueWorker(
 11032        IBookmarkQueueSignaler signaler,
 11033        IServiceScopeFactory scopeFactory,
 11034        ILogger<BookmarkQueueWorker> logger,
 11035        ITenantScopeFactory? tenantScopeFactory,
 11036        ITenantAccessor? tenantAccessor,
 11037        TimeSpan processThrottle)
 38    {
 11039        Signaler = signaler;
 11040        _scopeFactory = scopeFactory;
 11041        _tenantScopeFactory = tenantScopeFactory;
 11042        _tenantAccessor = tenantAccessor;
 11043        _logger = logger;
 11044        _tenant = tenantAccessor?.Tenant;
 11045        if (processThrottle > TimeSpan.Zero)
 8346            _rateLimitedProcessAsync = Throttler.Throttle<CancellationToken, Task>(ProcessAsync, processThrottle);
 11047    }
 48
 49    public void Start()
 50    {
 11051        if (_running)
 052            return;
 53
 11054        _tenant ??= _tenantAccessor?.Tenant;
 11055        _cts = new();
 11056        _running = true;
 57
 11058        _ = Task.Run(AwaitSignalAsync);
 11059    }
 60
 61    public void Stop()
 62    {
 3163        if (_running)
 64        {
 3165            _running = false;
 3166            _cts.Cancel();
 67            // Release is on the concrete type so IBookmarkQueueSignaler stays unchanged; a decorated signaler is left a
 3168            if (Signaler is BookmarkQueueSignaler bookmarkQueueSignaler)
 3169                bookmarkQueueSignaler.Release(CapturedTenantId);
 70        }
 71
 3172        _cts.Dispose();
 3173    }
 74
 75    private async Task AwaitSignalAsync()
 76    {
 49077        while (!_cts.IsCancellationRequested)
 78        {
 79            try
 80            {
 48781                using var tenantContext = _tenantAccessor?.PushContext(_tenant);
 48782                await Signaler.AwaitAsync(_cts.Token);
 38183                await InvokeProcessAsync(_cts.Token);
 38084            }
 2885            catch (OperationCanceledException)
 86            {
 2887                break; // Stop() was called
 88            }
 089            catch (Exception ex)
 90            {
 091                _logger.LogError(ex, "BookmarkQueueWorker error – continuing loop");
 092            }
 93        }
 3194    }
 95
 96    private Task InvokeProcessAsync(CancellationToken cancellationToken)
 97    {
 38198        return _rateLimitedProcessAsync is not null
 38199            ? _rateLimitedProcessAsync.InvokeAsync(cancellationToken)
 381100            : ProcessAsync(cancellationToken);
 101    }
 102
 103    protected virtual async Task ProcessAsync(CancellationToken cancellationToken)
 104    {
 372105        _logger.LogDebug("Processing bookmark queue...");
 106
 372107        if (_tenantScopeFactory is not null && _tenant is not null)
 108        {
 371109            await using var tenantScope = _tenantScopeFactory.CreateScope(_tenant);
 371110            var processor = tenantScope.ServiceProvider.GetRequiredService<IBookmarkQueueProcessor>();
 371111            await processor.ProcessAsync(cancellationToken);
 371112        }
 113        else
 114        {
 1115            using var scope = _scopeFactory.CreateScope();
 1116            var processor = scope.ServiceProvider.GetRequiredService<IBookmarkQueueProcessor>();
 1117            await processor.ProcessAsync(cancellationToken);
 1118        }
 119
 372120        _logger.LogDebug("Processed bookmark queue.");
 372121    }
 122}