< Summary - CoreWCF Coverage — PR #1766

Information
Class: CoreWCF.Runtime.InputQueue<T>
Assembly: CoreWCF.Primitives
File(s): /home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.Primitives/src/CoreWCF/Runtime/InputQueue.cs
Line coverage
0%
Covered lines: 0
Uncovered lines: 327
Coverable lines: 327
Total lines: 842
Line coverage: 0%
Branch coverage
0%
Covered branches: 0
Total branches: 158
Branch coverage: 0%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

File(s)

/home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.Primitives/src/CoreWCF/Runtime/InputQueue.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.Runtime
 10{
 11    internal sealed class InputQueue<T> : IDisposable where T : class
 12    {
 13        private static Action<object> s_completeOutstandingReadersCallback;
 14        private static Action<object> s_completeWaitersFalseCallback;
 15        private static Action<object> s_completeWaitersTrueCallback;
 16        private static Action<object> s_onDispatchCallback;
 17        private static Action<object> s_onInvokeDequeuedCallback;
 18        private QueueState _queueState;
 19        private readonly ItemQueue _itemQueue;
 20        private readonly Queue<IQueueReader> _readerQueue;
 21        private readonly List<IQueueWaiter> _waiterList;
 22
 023        public InputQueue()
 24        {
 025            _itemQueue = new ItemQueue();
 026            _readerQueue = new Queue<IQueueReader>();
 027            _waiterList = new List<IQueueWaiter>();
 028            _queueState = QueueState.Open;
 029        }
 30
 31        public InputQueue(Func<Action<AsyncCallback, IAsyncResult>> asyncCallbackGenerator)
 032            : this()
 33        {
 34            Fx.Assert(asyncCallbackGenerator != null, "use default ctor if you don't have a generator");
 035            AsyncCallbackGenerator = asyncCallbackGenerator;
 036        }
 37
 38        public int PendingCount
 39        {
 40            get
 41            {
 042                lock (ThisLock)
 43                {
 044                    return _itemQueue.ItemCount;
 45                }
 046            }
 47        }
 48
 49        // Users like ServiceModel can hook this abort ICommunicationObject or handle other non-IDisposable objects
 50        public Action<T> DisposeItemCallback
 51        {
 052            get;
 053            set;
 54        }
 55
 56        // Users like ServiceModel can hook this to wrap the AsyncQueueReader callback functionality for tracing, etc
 57        private Func<Action<AsyncCallback, IAsyncResult>> AsyncCallbackGenerator
 58        {
 059            get;
 060            set;
 61        }
 62
 63        private object ThisLock
 64        {
 065            get { return _itemQueue; }
 66        }
 67
 68        public void Close()
 69        {
 070            Dispose();
 071        }
 72
 73        public async Task<T> DequeueAsync(CancellationToken token)
 74        {
 075            (T result, bool success) dequeued = await TryDequeueAsync(token);
 76
 077            if (!dequeued.success)
 78            {
 79                // TODO: Create derived CancellationToken which carries original timeout with it
 080                throw Fx.Exception.AsError(new TimeoutException(SR.Format(SR.TimeoutInputQueueDequeue, null)));
 81            }
 82
 083            return dequeued.result;
 084        }
 85
 86        public async Task<(T result, bool success)> TryDequeueAsync(CancellationToken token)
 87        {
 088            WaitQueueReader reader = null;
 089            Item item = new Item();
 90
 091            lock (ThisLock)
 92            {
 093                if (_queueState == QueueState.Open)
 94                {
 095                    if (_itemQueue.HasAvailableItem)
 96                    {
 097                        item = _itemQueue.DequeueAvailableItem();
 98                    }
 99                    else
 100                    {
 0101                        reader = new WaitQueueReader(this);
 0102                        _readerQueue.Enqueue(reader);
 103                    }
 104                }
 0105                else if (_queueState == QueueState.Shutdown)
 106                {
 0107                    if (_itemQueue.HasAvailableItem)
 108                    {
 0109                        item = _itemQueue.DequeueAvailableItem();
 110                    }
 0111                    else if (_itemQueue.HasAnyItem)
 112                    {
 0113                        reader = new WaitQueueReader(this);
 0114                        _readerQueue.Enqueue(reader);
 115                    }
 116                    else
 117                    {
 0118                        return (default(T), true);
 119                    }
 120                }
 121                else // queueState == QueueState.Closed
 122                {
 0123                    return (default(T), true);
 124                }
 0125            }
 126
 0127            if (reader != null)
 128            {
 0129                return await reader.WaitAsync(token);
 130            }
 131            else
 132            {
 0133                InvokeDequeuedCallback(item.DequeuedCallback);
 0134                return (item.GetValue(), true);
 135            }
 0136        }
 137
 138        public void Dispatch()
 139        {
 0140            IQueueReader reader = null;
 0141            Item item = new Item();
 0142            IQueueReader[] outstandingReaders = null;
 0143            IQueueWaiter[] waiters = null;
 0144            bool itemAvailable = true;
 145
 0146            lock (ThisLock)
 147            {
 0148                itemAvailable = !((_queueState == QueueState.Closed) || (_queueState == QueueState.Shutdown));
 0149                GetWaiters(out waiters);
 150
 0151                if (_queueState != QueueState.Closed)
 152                {
 0153                    _itemQueue.MakePendingItemAvailable();
 154
 0155                    if (_readerQueue.Count > 0)
 156                    {
 0157                        item = _itemQueue.DequeueAvailableItem();
 0158                        reader = _readerQueue.Dequeue();
 159
 0160                        if (_queueState == QueueState.Shutdown && _readerQueue.Count > 0 && _itemQueue.ItemCount == 0)
 161                        {
 0162                            outstandingReaders = new IQueueReader[_readerQueue.Count];
 0163                            _readerQueue.CopyTo(outstandingReaders, 0);
 0164                            _readerQueue.Clear();
 165
 0166                            itemAvailable = false;
 167                        }
 168                    }
 169                }
 0170            }
 171
 0172            if (outstandingReaders != null)
 173            {
 0174                if (s_completeOutstandingReadersCallback == null)
 175                {
 0176                    s_completeOutstandingReadersCallback = CompleteOutstandingReadersCallback;
 177                }
 178
 0179                ActionItem.Schedule(s_completeOutstandingReadersCallback, outstandingReaders);
 180            }
 181
 0182            if (waiters != null)
 183            {
 0184                CompleteWaitersLater(itemAvailable, waiters);
 185            }
 186
 0187            if (reader != null)
 188            {
 0189                InvokeDequeuedCallback(item.DequeuedCallback);
 0190                reader.Set(item);
 191            }
 0192        }
 193
 194        public void EnqueueAndDispatch(T item)
 195        {
 0196            EnqueueAndDispatch(item, null);
 0197        }
 198
 199        // dequeuedCallback is called as an item is dequeued from the InputQueue.  The
 200        // InputQueue lock is not held during the callback.  However, the user code will
 201        // not be notified of the item being available until the callback returns.  If you
 202        // are not sure if the callback will block for a long time, then first call
 203        // IOThreadScheduler.ScheduleCallback to get to a "safe" thread.
 204        public void EnqueueAndDispatch(T item, Action dequeuedCallback)
 205        {
 0206            EnqueueAndDispatch(item, dequeuedCallback, true);
 0207        }
 208
 209        public void EnqueueAndDispatch(Exception exception, Action dequeuedCallback, bool canDispatchOnThisThread)
 210        {
 211            Fx.Assert(exception != null, "EnqueueAndDispatch: exception parameter should not be null");
 0212            EnqueueAndDispatch(new Item(exception, dequeuedCallback), canDispatchOnThisThread);
 0213        }
 214
 215        public void EnqueueAndDispatch(T item, Action dequeuedCallback, bool canDispatchOnThisThread)
 216        {
 217            Fx.Assert(item != null, "EnqueueAndDispatch: item parameter should not be null");
 0218            EnqueueAndDispatch(new Item(item, dequeuedCallback), canDispatchOnThisThread);
 0219        }
 220
 221        public bool EnqueueWithoutDispatch(T item, Action dequeuedCallback)
 222        {
 223            Fx.Assert(item != null, "EnqueueWithoutDispatch: item parameter should not be null");
 0224            return EnqueueWithoutDispatch(new Item(item, dequeuedCallback));
 225        }
 226
 227        public bool EnqueueWithoutDispatch(Exception exception, Action dequeuedCallback)
 228        {
 229            Fx.Assert(exception != null, "EnqueueWithoutDispatch: exception parameter should not be null");
 0230            return EnqueueWithoutDispatch(new Item(exception, dequeuedCallback));
 231        }
 232
 233
 234        public void Shutdown()
 235        {
 0236            Shutdown(null);
 0237        }
 238
 239        // Don't let any more items in. Differs from Close in that we keep around
 240        // existing items in our itemQueue for possible future calls to Dequeue
 241        public void Shutdown(Func<Exception> pendingExceptionGenerator)
 242        {
 0243            IQueueReader[] outstandingReaders = null;
 0244            lock (ThisLock)
 245            {
 0246                if (_queueState == QueueState.Shutdown)
 247                {
 0248                    return;
 249                }
 250
 0251                if (_queueState == QueueState.Closed)
 252                {
 0253                    return;
 254                }
 255
 0256                _queueState = QueueState.Shutdown;
 257
 0258                if (_readerQueue.Count > 0 && _itemQueue.ItemCount == 0)
 259                {
 0260                    outstandingReaders = new IQueueReader[_readerQueue.Count];
 0261                    _readerQueue.CopyTo(outstandingReaders, 0);
 0262                    _readerQueue.Clear();
 263                }
 0264            }
 265
 0266            if (outstandingReaders != null)
 267            {
 0268                for (int i = 0; i < outstandingReaders.Length; i++)
 269                {
 0270                    Exception exception = (pendingExceptionGenerator != null) ? pendingExceptionGenerator() : null;
 0271                    outstandingReaders[i].Set(new Item(exception, null));
 272                }
 273            }
 0274        }
 275
 276
 277        public Task<bool> WaitForItemAsync(CancellationToken token)
 278        {
 0279            WaitQueueWaiter waiter = null;
 0280            bool itemAvailable = false;
 281
 0282            lock (ThisLock)
 283            {
 0284                if (_queueState == QueueState.Open)
 285                {
 0286                    if (_itemQueue.HasAvailableItem)
 287                    {
 0288                        itemAvailable = true;
 289                    }
 290                    else
 291                    {
 0292                        waiter = new WaitQueueWaiter();
 0293                        _waiterList.Add(waiter);
 294                    }
 295                }
 0296                else if (_queueState == QueueState.Shutdown)
 297                {
 0298                    if (_itemQueue.HasAvailableItem)
 299                    {
 0300                        itemAvailable = true;
 301                    }
 0302                    else if (_itemQueue.HasAnyItem)
 303                    {
 0304                        waiter = new WaitQueueWaiter();
 0305                        _waiterList.Add(waiter);
 306                    }
 307                    else
 308                    {
 0309                        return Task.FromResult(true);
 310                    }
 311                }
 312                else // queueState == QueueState.Closed
 313                {
 0314                    return Task.FromResult(true);
 315                }
 316            }
 317
 0318            if (waiter != null)
 319            {
 0320                return waiter.WaitAsync(token);
 321            }
 322            else
 323            {
 0324                return Task.FromResult(itemAvailable);
 325            }
 0326        }
 327
 328        public void Dispose()
 329        {
 0330            bool dispose = false;
 331
 0332            lock (ThisLock)
 333            {
 0334                if (_queueState != QueueState.Closed)
 335                {
 0336                    _queueState = QueueState.Closed;
 0337                    dispose = true;
 338                }
 0339            }
 340
 0341            if (dispose)
 342            {
 0343                while (_readerQueue.Count > 0)
 344                {
 0345                    IQueueReader reader = _readerQueue.Dequeue();
 0346                    reader.Set(default);
 347                }
 348
 0349                while (_itemQueue.HasAnyItem)
 350                {
 0351                    Item item = _itemQueue.DequeueAnyItem();
 0352                    DisposeItem(item);
 0353                    InvokeDequeuedCallback(item.DequeuedCallback);
 354                }
 355            }
 0356        }
 357
 358        private void DisposeItem(Item item)
 359        {
 0360            T value = item.Value;
 0361            if (value != null)
 362            {
 0363                if (value is IDisposable)
 364                {
 0365                    ((IDisposable)value).Dispose();
 366                }
 367                else
 368                {
 0369                    Action<T> disposeItemCallback = DisposeItemCallback;
 0370                    if (disposeItemCallback != null)
 371                    {
 0372                        disposeItemCallback(value);
 373                    }
 374                }
 375            }
 0376        }
 377
 378        private static void CompleteOutstandingReadersCallback(object state)
 379        {
 0380            IQueueReader[] outstandingReaders = (IQueueReader[])state;
 381
 0382            for (int i = 0; i < outstandingReaders.Length; i++)
 383            {
 0384                outstandingReaders[i].Set(default);
 385            }
 0386        }
 387
 388        private static void CompleteWaiters(bool itemAvailable, IQueueWaiter[] waiters)
 389        {
 0390            for (int i = 0; i < waiters.Length; i++)
 391            {
 0392                waiters[i].Set(itemAvailable);
 393            }
 0394        }
 395
 396        private static void CompleteWaitersFalseCallback(object state)
 397        {
 0398            CompleteWaiters(false, (IQueueWaiter[])state);
 0399        }
 400
 401        private static void CompleteWaitersLater(bool itemAvailable, IQueueWaiter[] waiters)
 402        {
 0403            if (itemAvailable)
 404            {
 0405                if (s_completeWaitersTrueCallback == null)
 406                {
 0407                    s_completeWaitersTrueCallback = CompleteWaitersTrueCallback;
 408                }
 409
 0410                ActionItem.Schedule(s_completeWaitersTrueCallback, waiters);
 411            }
 412            else
 413            {
 0414                if (s_completeWaitersFalseCallback == null)
 415                {
 0416                    s_completeWaitersFalseCallback = CompleteWaitersFalseCallback;
 417                }
 418
 0419                ActionItem.Schedule(s_completeWaitersFalseCallback, waiters);
 420            }
 0421        }
 422
 423        private static void CompleteWaitersTrueCallback(object state)
 424        {
 0425            CompleteWaiters(true, (IQueueWaiter[])state);
 0426        }
 427
 428        private static void InvokeDequeuedCallback(Action dequeuedCallback)
 429        {
 0430            if (dequeuedCallback != null)
 431            {
 0432                dequeuedCallback();
 433            }
 0434        }
 435
 436        private static void InvokeDequeuedCallbackLater(Action dequeuedCallback)
 437        {
 0438            if (dequeuedCallback != null)
 439            {
 0440                if (s_onInvokeDequeuedCallback == null)
 441                {
 0442                    s_onInvokeDequeuedCallback = OnInvokeDequeuedCallback;
 443                }
 444
 0445                ActionItem.Schedule(s_onInvokeDequeuedCallback, dequeuedCallback);
 446            }
 0447        }
 448
 449        private static void OnDispatchCallback(object state)
 450        {
 0451            ((InputQueue<T>)state).Dispatch();
 0452        }
 453
 454        private static void OnInvokeDequeuedCallback(object state)
 455        {
 456            Fx.Assert(state != null, "InputQueue.OnInvokeDequeuedCallback: (state != null)");
 457
 0458            Action dequeuedCallback = (Action)state;
 0459            dequeuedCallback();
 0460        }
 461
 462        private void EnqueueAndDispatch(Item item, bool canDispatchOnThisThread)
 463        {
 0464            bool disposeItem = false;
 0465            IQueueReader reader = null;
 0466            bool dispatchLater = false;
 0467            IQueueWaiter[] waiters = null;
 0468            bool itemAvailable = true;
 469
 0470            lock (ThisLock)
 471            {
 0472                itemAvailable = !((_queueState == QueueState.Closed) || (_queueState == QueueState.Shutdown));
 0473                GetWaiters(out waiters);
 474
 0475                if (_queueState == QueueState.Open)
 476                {
 0477                    if (canDispatchOnThisThread)
 478                    {
 0479                        if (_readerQueue.Count == 0)
 480                        {
 0481                            _itemQueue.EnqueueAvailableItem(item);
 482                        }
 483                        else
 484                        {
 0485                            reader = _readerQueue.Dequeue();
 486                        }
 487                    }
 488                    else
 489                    {
 0490                        if (_readerQueue.Count == 0)
 491                        {
 0492                            _itemQueue.EnqueueAvailableItem(item);
 493                        }
 494                        else
 495                        {
 0496                            _itemQueue.EnqueuePendingItem(item);
 0497                            dispatchLater = true;
 498                        }
 499                    }
 500                }
 501                else // queueState == QueueState.Closed || queueState == QueueState.Shutdown
 502                {
 0503                    disposeItem = true;
 504                }
 0505            }
 506
 0507            if (waiters != null)
 508            {
 0509                if (canDispatchOnThisThread)
 510                {
 0511                    CompleteWaiters(itemAvailable, waiters);
 512                }
 513                else
 514                {
 0515                    CompleteWaitersLater(itemAvailable, waiters);
 516                }
 517            }
 518
 0519            if (reader != null)
 520            {
 0521                InvokeDequeuedCallback(item.DequeuedCallback);
 0522                reader.Set(item);
 523            }
 524
 0525            if (dispatchLater)
 526            {
 0527                if (s_onDispatchCallback == null)
 528                {
 0529                    s_onDispatchCallback = OnDispatchCallback;
 530                }
 531
 0532                ActionItem.Schedule(s_onDispatchCallback, this);
 533            }
 0534            else if (disposeItem)
 535            {
 0536                InvokeDequeuedCallback(item.DequeuedCallback);
 0537                DisposeItem(item);
 538            }
 0539        }
 540
 541        // This will not block, however, Dispatch() must be called later if this function
 542        // returns true.
 543        private bool EnqueueWithoutDispatch(Item item)
 544        {
 0545            lock (ThisLock)
 546            {
 547                // Open
 0548                if (_queueState != QueueState.Closed && _queueState != QueueState.Shutdown)
 549                {
 0550                    if (_readerQueue.Count == 0 && _waiterList.Count == 0)
 551                    {
 0552                        _itemQueue.EnqueueAvailableItem(item);
 0553                        return false;
 554                    }
 555                    else
 556                    {
 0557                        _itemQueue.EnqueuePendingItem(item);
 0558                        return true;
 559                    }
 560                }
 0561            }
 562
 0563            DisposeItem(item);
 0564            InvokeDequeuedCallbackLater(item.DequeuedCallback);
 0565            return false;
 0566        }
 567
 568        private void GetWaiters(out IQueueWaiter[] waiters)
 569        {
 0570            if (_waiterList.Count > 0)
 571            {
 0572                waiters = _waiterList.ToArray();
 0573                _waiterList.Clear();
 574            }
 575            else
 576            {
 0577                waiters = null;
 578            }
 0579        }
 580
 581        // Used for timeouts. The InputQueue must remove readers from its reader queue to prevent
 582        // dispatching items to timed out readers.
 583        private bool RemoveReader(IQueueReader reader)
 584        {
 585            Fx.Assert(reader != null, "InputQueue.RemoveReader: (reader != null)");
 586
 0587            lock (ThisLock)
 588            {
 0589                if (_queueState == QueueState.Open || _queueState == QueueState.Shutdown)
 590                {
 0591                    bool removed = false;
 592
 0593                    for (int i = _readerQueue.Count; i > 0; i--)
 594                    {
 0595                        IQueueReader temp = _readerQueue.Dequeue();
 0596                        if (ReferenceEquals(temp, reader))
 597                        {
 0598                            removed = true;
 599                        }
 600                        else
 601                        {
 0602                            _readerQueue.Enqueue(temp);
 603                        }
 604                    }
 605
 0606                    return removed;
 607                }
 0608            }
 609
 0610            return false;
 0611        }
 612
 613        private enum QueueState
 614        {
 615            Open,
 616            Shutdown,
 617            Closed
 618        }
 619
 620        private interface IQueueReader
 621        {
 622            void Set(Item item);
 623        }
 624
 625        private interface IQueueWaiter
 626        {
 627            void Set(bool itemAvailable);
 628        }
 629
 630        private struct Item
 631        {
 632            public Item(T value, Action dequeuedCallback)
 0633                : this(value, null, dequeuedCallback)
 634            {
 0635            }
 636
 637            public Item(Exception exception, Action dequeuedCallback)
 0638                : this(null, exception, dequeuedCallback)
 639            {
 0640            }
 641
 642            private Item(T value, Exception exception, Action dequeuedCallback)
 643            {
 0644                Value = value;
 0645                Exception = exception;
 0646                DequeuedCallback = dequeuedCallback;
 0647            }
 648
 0649            public Action DequeuedCallback { get; }
 650
 0651            public Exception Exception { get; }
 652
 0653            public T Value { get; }
 654
 655            public T GetValue()
 656            {
 0657                if (Exception != null)
 658                {
 0659                    throw Fx.Exception.AsError(Exception);
 660                }
 661
 0662                return Value;
 663            }
 664        }
 665
 666        private class ItemQueue
 667        {
 668            private int _head;
 669            private Item[] _items;
 670            private int _pendingCount;
 671
 0672            public ItemQueue()
 673            {
 0674                _items = new Item[1];
 0675            }
 676
 677            public bool HasAnyItem
 678            {
 0679                get { return ItemCount > 0; }
 680            }
 681
 682            public bool HasAvailableItem
 683            {
 0684                get { return ItemCount > _pendingCount; }
 685            }
 686
 0687            public int ItemCount { get; private set; }
 688
 689            public Item DequeueAnyItem()
 690            {
 0691                if (_pendingCount == ItemCount)
 692                {
 0693                    _pendingCount--;
 694                }
 0695                return DequeueItemCore();
 696            }
 697
 698            public Item DequeueAvailableItem()
 699            {
 0700                Fx.AssertAndThrow(ItemCount != _pendingCount, "ItemQueue does not contain any available items");
 0701                return DequeueItemCore();
 702            }
 703
 704            public void EnqueueAvailableItem(Item item)
 705            {
 0706                EnqueueItemCore(item);
 0707            }
 708
 709            public void EnqueuePendingItem(Item item)
 710            {
 0711                EnqueueItemCore(item);
 0712                _pendingCount++;
 0713            }
 714
 715            public void MakePendingItemAvailable()
 716            {
 0717                Fx.AssertAndThrow(_pendingCount != 0, "ItemQueue does not contain any pending items");
 0718                _pendingCount--;
 0719            }
 720
 721            private Item DequeueItemCore()
 722            {
 0723                Fx.AssertAndThrow(ItemCount != 0, "ItemQueue does not contain any items");
 0724                Item item = _items[_head];
 0725                _items[_head] = new Item();
 0726                ItemCount--;
 0727                _head = (_head + 1) % _items.Length;
 0728                return item;
 729            }
 730
 731            private void EnqueueItemCore(Item item)
 732            {
 0733                if (ItemCount == _items.Length)
 734                {
 0735                    Item[] newItems = new Item[_items.Length * 2];
 0736                    for (int i = 0; i < ItemCount; i++)
 737                    {
 0738                        newItems[i] = _items[(_head + i) % _items.Length];
 739                    }
 0740                    _head = 0;
 0741                    _items = newItems;
 742                }
 0743                int tail = (_head + ItemCount) % _items.Length;
 0744                _items[tail] = item;
 0745                ItemCount++;
 0746            }
 747        }
 748
 749        private class WaitQueueReader : IQueueReader
 750        {
 751            private Exception _exception;
 752            private readonly InputQueue<T> _inputQueue;
 753            private T _item;
 754            private readonly AsyncManualResetEvent _waitEvent;
 755
 0756            public WaitQueueReader(InputQueue<T> inputQueue)
 757            {
 0758                _inputQueue = inputQueue;
 0759                _waitEvent = new AsyncManualResetEvent();
 0760            }
 761
 762            public void Set(Item item)
 763            {
 0764                lock (this)
 765                {
 766                    Fx.Assert(_item == null, "InputQueue.WaitQueueReader.Set: (this.item == null)");
 767                    Fx.Assert(_exception == null, "InputQueue.WaitQueueReader.Set: (this.exception == null)");
 768
 0769                    _exception = item.Exception;
 0770                    _item = item.Value;
 0771                    _waitEvent.Set();
 0772                }
 0773            }
 774
 775            public async Task<(T result, bool success)> WaitAsync(CancellationToken token)
 776            {
 0777                bool isSafeToClose = false;
 778                try
 779                {
 0780                    if (!await _waitEvent.WaitAsync(token))
 781                    {
 0782                        if (_inputQueue.RemoveReader(this))
 783                        {
 0784                            isSafeToClose = true;
 0785                            return (null, false);
 786                        }
 787                        else
 788                        {
 0789                            await _waitEvent.WaitAsync();
 790                        }
 791                    }
 792
 0793                    isSafeToClose = true;
 0794                }
 795                finally
 796                {
 0797                    if (isSafeToClose)
 798                    {
 0799                        _waitEvent.Dispose();
 800                    }
 801                }
 802
 0803                if (_exception != null)
 804                {
 0805                    throw Fx.Exception.AsError(_exception);
 806                }
 807
 0808                return (_item, true);
 0809            }
 810        }
 811
 812        private class WaitQueueWaiter : IQueueWaiter
 813        {
 814            private bool _itemAvailable;
 815            private readonly AsyncManualResetEvent _waitEvent;
 816
 0817            public WaitQueueWaiter()
 818            {
 0819                _waitEvent = new AsyncManualResetEvent();
 0820            }
 821
 822            public void Set(bool itemAvailable)
 823            {
 0824                lock (this)
 825                {
 0826                    _itemAvailable = itemAvailable;
 0827                    _waitEvent.Set();
 0828                }
 0829            }
 830
 831            public async Task<bool> WaitAsync(CancellationToken token)
 832            {
 0833                if (!await _waitEvent.WaitAsync(token))
 834                {
 0835                    return false;
 836                }
 837
 0838                return _itemAvailable;
 0839            }
 840        }
 841    }
 842}

Methods/Properties

.ctor()
.ctor(System.Func`1<System.Action`2<System.AsyncCallback,System.IAsyncResult>>)
PendingCount()
DisposeItemCallback()
DisposeItemCallback(System.Action`1<T>)
AsyncCallbackGenerator()
AsyncCallbackGenerator(System.Func`1<System.Action`2<System.AsyncCallback,System.IAsyncResult>>)
ThisLock()
Close()
DequeueAsync()
TryDequeueAsync()
Dispatch()
EnqueueAndDispatch(T)
EnqueueAndDispatch(T,System.Action)
EnqueueAndDispatch(System.Exception,System.Action,System.Boolean)
EnqueueAndDispatch(T,System.Action,System.Boolean)
EnqueueWithoutDispatch(T,System.Action)
EnqueueWithoutDispatch(System.Exception,System.Action)
Shutdown()
Shutdown(System.Func`1<System.Exception>)
WaitForItemAsync(System.Threading.CancellationToken)
Dispose()
DisposeItem(CoreWCF.Runtime.InputQueue`1/Item<T>)
CompleteOutstandingReadersCallback(System.Object)
CompleteWaiters(System.Boolean,CoreWCF.Runtime.InputQueue`1/IQueueWaiter<T>[])
CompleteWaitersFalseCallback(System.Object)
CompleteWaitersLater(System.Boolean,CoreWCF.Runtime.InputQueue`1/IQueueWaiter<T>[])
CompleteWaitersTrueCallback(System.Object)
InvokeDequeuedCallback(System.Action)
InvokeDequeuedCallbackLater(System.Action)
OnDispatchCallback(System.Object)
OnInvokeDequeuedCallback(System.Object)
EnqueueAndDispatch(CoreWCF.Runtime.InputQueue`1/Item<T>,System.Boolean)
EnqueueWithoutDispatch(CoreWCF.Runtime.InputQueue`1/Item<T>)
GetWaiters(CoreWCF.Runtime.InputQueue`1/IQueueWaiter<T>[]&)
RemoveReader(CoreWCF.Runtime.InputQueue`1/IQueueReader<T>)
.ctor(T,System.Action)
.ctor(System.Exception,System.Action)
.ctor(T,System.Exception,System.Action)
DequeuedCallback()
Exception()
Value()
GetValue()
.ctor()
HasAnyItem()
HasAvailableItem()
ItemCount()
DequeueAnyItem()
DequeueAvailableItem()
EnqueueAvailableItem(CoreWCF.Runtime.InputQueue`1/Item<T>)
EnqueuePendingItem(CoreWCF.Runtime.InputQueue`1/Item<T>)
MakePendingItemAvailable()
DequeueItemCore()
EnqueueItemCore(CoreWCF.Runtime.InputQueue`1/Item<T>)
.ctor(CoreWCF.Runtime.InputQueue`1<T>)
Set(CoreWCF.Runtime.InputQueue`1/Item<T>)
WaitAsync()
.ctor()
Set(System.Boolean)
WaitAsync()