< Summary

Information
Class: Elsa.Workflows.Runtime.BookmarkQueueProcessor
Assembly: Elsa.Workflows.Runtime
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/BookmarkQueueProcessor.cs
Line coverage
100%
Covered lines: 61
Uncovered lines: 0
Coverable lines: 61
Total lines: 116
Line coverage: 100%
Branch coverage
100%
Covered branches: 12
Total branches: 12
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
ProcessAsync()100%44100%
ProcessPageAsync()100%44100%
ProcessItemAsync()100%2287.5%
HandleFailureAsync()100%22100%

File(s)

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

#LineLine coverage
 1using Elsa.Common.Entities;
 2using Elsa.Common.Models;
 3using Elsa.Extensions;
 4using Elsa.Common;
 5using Elsa.Workflows.Runtime.Entities;
 6using Elsa.Workflows.Runtime.Messages;
 7using Elsa.Workflows.Runtime.Options;
 8using Elsa.Workflows.Runtime.OrderDefinitions;
 9using Microsoft.Extensions.Logging;
 10using Microsoft.Extensions.Options;
 11
 12namespace Elsa.Workflows.Runtime;
 13
 33514public class BookmarkQueueProcessor(
 33515    IBookmarkQueueStore store,
 33516    IBookmarkQueueDeadLetterManager deadLetterManager,
 33517    IWorkflowResumer workflowResumer,
 33518    ISystemClock systemClock,
 33519    IOptions<BookmarkQueuePurgeOptions> options,
 33520    ILogger<BookmarkQueueProcessor> logger) : IBookmarkQueueProcessor
 21{
 22    public async Task ProcessAsync(CancellationToken cancellationToken = default)
 23    {
 33624        var batchSize = 50;
 33625        var offset = 0;
 26
 33727        while (!cancellationToken.IsCancellationRequested)
 28        {
 33629            var pageArgs = PageArgs.FromRange(offset, batchSize);
 33630            var page = await store.PageAsync(pageArgs, new BookmarkQueueItemOrder<DateTimeOffset>(x => x.CreatedAt, Orde
 33631            var retainedCount = await ProcessPageAsync(page, cancellationToken);
 32
 33533            if (page.Items.Count < batchSize)
 34                break;
 35
 136            offset += retainedCount;
 137        }
 33538    }
 39
 40    private async Task<int> ProcessPageAsync(Page<BookmarkQueueItem> page, CancellationToken cancellationToken = default
 41    {
 33642        var retainedCount = 0;
 43
 1363344        foreach (var bookmarkQueueItem in page.Items)
 45        {
 648146            if (await ProcessItemAsync(bookmarkQueueItem, cancellationToken))
 638147                retainedCount++;
 48        }
 49
 33550        return retainedCount;
 33551    }
 52
 53    private async Task<bool> ProcessItemAsync(BookmarkQueueItem item, CancellationToken cancellationToken = default)
 54    {
 648155        var filter = item.CreateBookmarkFilter();
 648156        var resumeOptions = item.Options;
 57
 648158        logger.LogDebug("Processing bookmark queue item {BookmarkQueueItemId} for workflow instance {WorkflowInstanceId}
 59
 60        List<RunWorkflowInstanceResponse> responses;
 61
 62        try
 63        {
 648164            responses = (await workflowResumer.ResumeAsync(filter, resumeOptions, cancellationToken)).ToList();
 647765        }
 166        catch (OperationCanceledException)
 67        {
 168            throw;
 69        }
 370        catch (Exception ex) when (ex is not OutOfMemoryException and not StackOverflowException)
 71        {
 372            return await HandleFailureAsync(item, ex, cancellationToken);
 73        }
 74
 647775        if (responses.Count > 0)
 76        {
 9777            logger.LogDebug("Successfully resumed {WorkflowCount} workflow instances using stimulus {StimulusHash} for a
 9778            await store.DeleteAsync(item.Id, cancellationToken);
 9779            return false;
 80        }
 81
 638082        logger.LogDebug("No matching bookmarks found for bookmark queue item {BookmarkQueueItemId} for workflow instance
 638083        return true;
 648084    }
 85
 86    private async Task<bool> HandleFailureAsync(BookmarkQueueItem item, Exception exception, CancellationToken cancellat
 87    {
 388        item.DeliveryAttempts++;
 389        item.LastAttemptedAt = systemClock.UtcNow;
 390        item.LastErrorType = exception.GetType().FullName;
 391        item.LastErrorMessage = exception.Message;
 92
 393        if (item.DeliveryAttempts < options.Value.MaxDeliveryAttempts)
 94        {
 195            logger.LogWarning(
 196                exception,
 197                "Failed to process bookmark queue item {BookmarkQueueItemId}. Attempt {DeliveryAttempt} of {MaxDeliveryA
 198                item.Id,
 199                item.DeliveryAttempts,
 1100                options.Value.MaxDeliveryAttempts);
 101
 1102            await store.SaveAsync(item, cancellationToken);
 1103            return true;
 104        }
 105
 2106        logger.LogError(
 2107            exception,
 2108            "Moving bookmark queue item {BookmarkQueueItemId} to dead letter after {DeliveryAttempt} failed delivery att
 2109            item.Id,
 2110            item.DeliveryAttempts);
 111
 2112        await deadLetterManager.DeadLetterAsync(item, "Failed", exception, cancellationToken);
 2113        await store.DeleteAsync(item.Id, cancellationToken);
 2114        return false;
 3115    }
 116}