| | | 1 | | using Elsa.Common.Multitenancy; |
| | | 2 | | using Microsoft.Extensions.DependencyInjection; |
| | | 3 | | using Microsoft.Extensions.Logging; |
| | | 4 | | using ThrottleDebounce; |
| | | 5 | | |
| | | 6 | | namespace Elsa.Workflows.Runtime; |
| | | 7 | | |
| | | 8 | | public 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; |
| | 519 | 18 | | protected IBookmarkQueueSignaler Signaler { get; } |
| | 41 | 19 | | 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) |
| | 81 | 27 | | : this(signaler, scopeFactory, logger, tenantScopeFactory, tenantAccessor, TimeSpan.FromMilliseconds(500)) |
| | | 28 | | { |
| | 81 | 29 | | } |
| | | 30 | | |
| | 110 | 31 | | protected BookmarkQueueWorker( |
| | 110 | 32 | | IBookmarkQueueSignaler signaler, |
| | 110 | 33 | | IServiceScopeFactory scopeFactory, |
| | 110 | 34 | | ILogger<BookmarkQueueWorker> logger, |
| | 110 | 35 | | ITenantScopeFactory? tenantScopeFactory, |
| | 110 | 36 | | ITenantAccessor? tenantAccessor, |
| | 110 | 37 | | TimeSpan processThrottle) |
| | | 38 | | { |
| | 110 | 39 | | Signaler = signaler; |
| | 110 | 40 | | _scopeFactory = scopeFactory; |
| | 110 | 41 | | _tenantScopeFactory = tenantScopeFactory; |
| | 110 | 42 | | _tenantAccessor = tenantAccessor; |
| | 110 | 43 | | _logger = logger; |
| | 110 | 44 | | _tenant = tenantAccessor?.Tenant; |
| | 110 | 45 | | if (processThrottle > TimeSpan.Zero) |
| | 83 | 46 | | _rateLimitedProcessAsync = Throttler.Throttle<CancellationToken, Task>(ProcessAsync, processThrottle); |
| | 110 | 47 | | } |
| | | 48 | | |
| | | 49 | | public void Start() |
| | | 50 | | { |
| | 110 | 51 | | if (_running) |
| | 0 | 52 | | return; |
| | | 53 | | |
| | 110 | 54 | | _tenant ??= _tenantAccessor?.Tenant; |
| | 110 | 55 | | _cts = new(); |
| | 110 | 56 | | _running = true; |
| | | 57 | | |
| | 110 | 58 | | _ = Task.Run(AwaitSignalAsync); |
| | 110 | 59 | | } |
| | | 60 | | |
| | | 61 | | public void Stop() |
| | | 62 | | { |
| | 31 | 63 | | if (_running) |
| | | 64 | | { |
| | 31 | 65 | | _running = false; |
| | 31 | 66 | | _cts.Cancel(); |
| | | 67 | | // Release is on the concrete type so IBookmarkQueueSignaler stays unchanged; a decorated signaler is left a |
| | 31 | 68 | | if (Signaler is BookmarkQueueSignaler bookmarkQueueSignaler) |
| | 31 | 69 | | bookmarkQueueSignaler.Release(CapturedTenantId); |
| | | 70 | | } |
| | | 71 | | |
| | 31 | 72 | | _cts.Dispose(); |
| | 31 | 73 | | } |
| | | 74 | | |
| | | 75 | | private async Task AwaitSignalAsync() |
| | | 76 | | { |
| | 490 | 77 | | while (!_cts.IsCancellationRequested) |
| | | 78 | | { |
| | | 79 | | try |
| | | 80 | | { |
| | 487 | 81 | | using var tenantContext = _tenantAccessor?.PushContext(_tenant); |
| | 487 | 82 | | await Signaler.AwaitAsync(_cts.Token); |
| | 381 | 83 | | await InvokeProcessAsync(_cts.Token); |
| | 380 | 84 | | } |
| | 28 | 85 | | catch (OperationCanceledException) |
| | | 86 | | { |
| | 28 | 87 | | break; // Stop() was called |
| | | 88 | | } |
| | 0 | 89 | | catch (Exception ex) |
| | | 90 | | { |
| | 0 | 91 | | _logger.LogError(ex, "BookmarkQueueWorker error – continuing loop"); |
| | 0 | 92 | | } |
| | | 93 | | } |
| | 31 | 94 | | } |
| | | 95 | | |
| | | 96 | | private Task InvokeProcessAsync(CancellationToken cancellationToken) |
| | | 97 | | { |
| | 381 | 98 | | return _rateLimitedProcessAsync is not null |
| | 381 | 99 | | ? _rateLimitedProcessAsync.InvokeAsync(cancellationToken) |
| | 381 | 100 | | : ProcessAsync(cancellationToken); |
| | | 101 | | } |
| | | 102 | | |
| | | 103 | | protected virtual async Task ProcessAsync(CancellationToken cancellationToken) |
| | | 104 | | { |
| | 372 | 105 | | _logger.LogDebug("Processing bookmark queue..."); |
| | | 106 | | |
| | 372 | 107 | | if (_tenantScopeFactory is not null && _tenant is not null) |
| | | 108 | | { |
| | 371 | 109 | | await using var tenantScope = _tenantScopeFactory.CreateScope(_tenant); |
| | 371 | 110 | | var processor = tenantScope.ServiceProvider.GetRequiredService<IBookmarkQueueProcessor>(); |
| | 371 | 111 | | await processor.ProcessAsync(cancellationToken); |
| | 371 | 112 | | } |
| | | 113 | | else |
| | | 114 | | { |
| | 1 | 115 | | using var scope = _scopeFactory.CreateScope(); |
| | 1 | 116 | | var processor = scope.ServiceProvider.GetRequiredService<IBookmarkQueueProcessor>(); |
| | 1 | 117 | | await processor.ProcessAsync(cancellationToken); |
| | 1 | 118 | | } |
| | | 119 | | |
| | 372 | 120 | | _logger.LogDebug("Processed bookmark queue."); |
| | 372 | 121 | | } |
| | | 122 | | } |