< Summary - CoreWCF Coverage — PR #1766

Information
Class: CoreWCF.Queue.Common.DefaultQueueTransportPump
Assembly: CoreWCF.Queue
File(s): /home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.Queue/src/CoreWCF/Queue/Common/DefaultQueueTransportPump.cs
Line coverage
96%
Covered lines: 29
Uncovered lines: 1
Coverable lines: 30
Total lines: 79
Line coverage: 96.6%
Branch coverage
87%
Covered branches: 7
Total branches: 8
Branch coverage: 87.5%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Cyclomatic complexity NPath complexity Sequence coverage
.ctor(...)100%11100%
StartPumpAsync(...)100%22100%
StopPumpAsync(...)50%2275%
FetchAndProcessAsync()100%44100%

File(s)

/home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.Queue/src/CoreWCF/Queue/Common/DefaultQueueTransportPump.cs

#LineLine coverage
 1// Licensed to the .NET Foundation under one or more agreements.
 2// The .NET Foundation licenses this file to you under the MIT license.
 3
 4using System;
 5using System.Collections.Generic;
 6using System.Threading;
 7using System.Threading.Tasks;
 8
 9namespace CoreWCF.Queue.Common
 10{
 11    /// <summary>
 12    /// DefaultQueueTransportPump model serves the purpose of pull model
 13    /// </summary>
 14    internal class DefaultQueueTransportPump : QueueTransportPump
 15    {
 16        private readonly IQueueTransport _transport;
 317        private readonly CancellationTokenSource _cancellationTokenSource = new();
 318        private readonly List<Task> _tasks = new();
 19
 320        public DefaultQueueTransportPump(IQueueTransport queueTransport)
 21        {
 322            _transport = queueTransport;
 323        }
 24
 25        public override Task StartPumpAsync(QueueTransportContext queueTransportContext, CancellationToken token)
 26        {
 1227            for (int i = 0; i < _transport.ConcurrencyLevel; i++)
 28            {
 329                _tasks.Add(FetchAndProcessAsync(queueTransportContext, _cancellationTokenSource.Token));
 30            }
 31
 332            return Task.CompletedTask;
 33        }
 34
 35        public override Task StopPumpAsync(CancellationToken token)
 36        {
 337            _cancellationTokenSource.Cancel();
 338            if (_transport is IDisposable disposable)
 39            {
 040                disposable.Dispose();
 41            }
 42
 343            return Task.WhenAll(_tasks);
 44        }
 45
 46        private async Task FetchAndProcessAsync(QueueTransportContext queueTransportContext, CancellationToken token)
 47        {
 348            CancellationTokenSource cts = new();
 349            TimeSpan receiveTimeout = queueTransportContext.ServiceDispatcher.Binding.ReceiveTimeout;
 30050            while (!token.IsCancellationRequested)
 51            {
 29752                cts.CancelAfter(receiveTimeout);
 53                QueueMessageContext queueMessageContext;
 54                try
 55                {
 29756                    using var linkedCts =
 29757                        CancellationTokenSource.CreateLinkedTokenSource(cts.Token, token);
 29758                    queueMessageContext = await _transport.ReceiveQueueMessageContextAsync(linkedCts.Token);
 19659                }
 10160                catch (OperationCanceledException)
 61                {
 10162                    cts = new();
 10163                    continue;
 64                }
 65
 19666                if (queueMessageContext == null)
 67                {
 9868                    cts = new();
 9869                    continue;
 70                }
 71
 9872                cts.CancelAfter(-1);
 73
 9874                queueMessageContext.QueueTransportContext = queueTransportContext;
 9875                await queueTransportContext.QueueMessageDispatcher(queueMessageContext);
 76            }
 377        }
 78    }
 79}