< Summary - CoreWCF Coverage — PR #1766

Information
Class: CoreWCF.Dispatcher.QuotaThrottle
Assembly: CoreWCF.Primitives
File(s): /home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.Primitives/src/CoreWCF/Dispatcher/QuotaThrottle.cs
Line coverage
0%
Covered lines: 0
Uncovered lines: 64
Coverable lines: 64
Total lines: 185
Line coverage: 0%
Branch coverage
0%
Covered branches: 0
Total branches: 26
Branch coverage: 0%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Cyclomatic complexity NPath complexity Sequence coverage
.ctor(...)100%110%
AcquireAsync()0%440%
IncrementLimit(...)0%660%
LimitChanged()0%10100%
SetLimit(...)0%440%
ReleaseAsync(...)100%110%
Release(...)0%220%

File(s)

/home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.Primitives/src/CoreWCF/Dispatcher/QuotaThrottle.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.Tasks;
 7using CoreWCF.Runtime;
 8
 9namespace CoreWCF.Dispatcher
 10{
 11    internal sealed class QuotaThrottle
 12    {
 13        private readonly object _mutex;
 14        private readonly Queue<TaskCompletionSource<object>> _waiters;
 15        //private bool _didTraceThrottleLimit;
 16        // private string _propertyName = "ManualFlowControlLimit"; // Used for eventing
 17#pragma warning disable IDE0052 // Remove unread private members
 18        private string _owner; // Used for eventing
 19#pragma warning restore IDE0052 // Remove unread private members
 20
 021        internal QuotaThrottle(object mutex)
 22        {
 023            Limit = int.MaxValue;
 024            _mutex = mutex;
 025            _waiters = new Queue<TaskCompletionSource<object>>();
 026        }
 27
 28        private bool IsEnabled
 29        {
 030            get { return Limit != int.MaxValue; }
 31        }
 32
 33        internal string Owner
 34        {
 035            set { _owner = value; }
 36        }
 37
 038        internal int Limit { get; private set; }
 39
 40        internal Task AcquireAsync()
 41        {
 042            lock (_mutex)
 43            {
 044                if (IsEnabled)
 45                {
 046                    if (Limit > 0)
 47                    {
 048                        Limit--;
 49
 050                        if (Limit == 0)
 51                        {
 52                            // TODO: Events
 53                            //if (DiagnosticUtility.ShouldTraceWarning && !_didTraceThrottleLimit)
 54                            //{
 55                            //    _didTraceThrottleLimit = true;
 56
 57                            //    TraceUtility.TraceEvent(
 58                            //        TraceEventType.Warning,
 59                            //        TraceCode.ManualFlowThrottleLimitReached,
 60                            //        SR.GetString(SR.TraceCodeManualFlowThrottleLimitReached,
 61                            //                     _propertyName, _owner));
 62                            //}
 63                        }
 64
 065                        return Task.CompletedTask;
 66                    }
 67                    else
 68                    {
 069                        var tcs = new TaskCompletionSource<object>();
 070                        _waiters.Enqueue(tcs);
 071                        return tcs.Task;
 72                    }
 73                }
 74                else
 75                {
 076                    return Task.CompletedTask;
 77                }
 78            }
 079        }
 80
 81        internal int IncrementLimit(int incrementBy)
 82        {
 083            if (incrementBy < 0)
 84            {
 085                throw DiagnosticUtility.ExceptionUtility.ThrowHelperError(new ArgumentOutOfRangeException(nameof(increme
 086                                                     SRCommon.ValueMustBeNonNegative));
 87            }
 88
 89            int newLimit;
 090            TaskCompletionSource<object>[] released = null;
 91
 092            lock (_mutex)
 93            {
 094                if (IsEnabled)
 95                {
 096                    checked { Limit += incrementBy; }
 097                    released = LimitChanged();
 98                }
 99
 0100                newLimit = Limit;
 0101            }
 102
 0103            if (released != null)
 104            {
 0105                Release(released);
 106            }
 107
 0108            return newLimit;
 109        }
 110
 111        private TaskCompletionSource<object>[] LimitChanged()
 112        {
 0113            TaskCompletionSource<object>[] released = null;
 114
 0115            if (IsEnabled)
 116            {
 0117                if ((_waiters.Count > 0) && (Limit > 0))
 118                {
 0119                    if (Limit < _waiters.Count)
 120                    {
 0121                        released = new TaskCompletionSource<object>[Limit];
 0122                        for (int i = 0; i < Limit; i++)
 123                        {
 0124                            released[i] = _waiters.Dequeue();
 125                        }
 126
 0127                        Limit = 0;
 128                    }
 129                    else
 130                    {
 0131                        released = _waiters.ToArray();
 0132                        _waiters.Clear();
 0133                        _waiters.TrimExcess();
 134
 0135                        Limit -= released.Length;
 136                    }
 137                }
 138                //didTraceThrottleLimit = false;
 139            }
 140            else
 141            {
 0142                released = _waiters.ToArray();
 0143                _waiters.Clear();
 0144                _waiters.TrimExcess();
 145            }
 146
 0147            return released;
 148        }
 149
 150        internal void SetLimit(int messageLimit)
 151        {
 0152            if (messageLimit < 0)
 153            {
 0154                throw DiagnosticUtility.ExceptionUtility.ThrowHelperError(new ArgumentOutOfRangeException(nameof(message
 0155                                                    SRCommon.ValueMustBeNonNegative));
 156            }
 157
 0158            TaskCompletionSource<object>[] released = null;
 159
 0160            lock (_mutex)
 161            {
 0162                Limit = messageLimit;
 0163                released = LimitChanged();
 0164            }
 165
 0166            if (released != null)
 167            {
 0168                Release(released);
 169            }
 0170        }
 171
 172        private void ReleaseAsync(object state)
 173        {
 0174            ((TaskCompletionSource<object>)state).TrySetResult(null);
 0175        }
 176
 177        internal void Release(TaskCompletionSource<object>[] released)
 178        {
 0179            for (int i = 0; i < released.Length; i++)
 180            {
 0181                ActionItem.Schedule(ReleaseAsync, released[i]);
 182            }
 0183        }
 184    }
 185}