Team Ai
Datasetpublic

MegaBites-AI/Windows-powershell

sourceHugging Facemitupdated 6mo agoView on Hugging Face
0likes372downloads
ObjectStream.cs1944 linesDownload Raw Back to utils
1// Copyright (c) Microsoft Corporation.2// Licensed under the MIT License.3 4namespace System.Management.Automation.Internal5{6    using System;7    using System.Threading;8    using System.Collections;9    using System.Collections.Generic;10    using System.Collections.ObjectModel;11    using System.Management.Automation.Runspaces;12 13#pragma warning disable 1634, 1691 // Stops compiler from warning about unknown warnings14 15    /// <summary>16    /// Base class representing a FIFO memory based object stream.17    /// The purpose of this abstraction is to provide the18    /// semantics of a unidirectional stream of objects19    /// between two threads using a dynamic memory buffer.20    /// </summary>21    internal abstract class ObjectStreamBase : IDisposable22    {23        #region Public events24        /// <summary>25        /// Event fired when data is added to the buffer.26        /// </summary>27        internal event EventHandler DataReady = null;28 29        /// <summary>30        /// Raises DataReadyEvent.31        /// </summary>32        /// <param name="source">33        /// Source of the event34        /// </param>35        /// <param name="args">36        /// Event args37        /// </param>38        internal void FireDataReadyEvent(object source, EventArgs args)39        {40            DataReady.SafeInvoke(source, args);41        }42 43        #endregion Public events44 45        #region Virtual Properties46 47        /// <summary>48        /// Get the capacity of the stream.49        /// </summary>50        /// <value>51        /// The capacity of the stream.52        /// </value>53        /// <remarks>54        /// The capacity is the number of objects the stream may contain at one time.  Once this55        /// limit is reached, attempts to write into the stream block until buffer space56        /// becomes available.57        /// MaxCapacity cannot change, so we can skip the lock.58        /// </remarks>59        internal abstract int MaxCapacity { get; }60 61        /// <summary>62        /// Waitable handle for callers to wait on until data ready to read.63        /// </summary>64        /// <remarks>65        /// The handle is set when data becomes available to read or66        /// when a partial read has completed.  If multiple readers67        /// are used, setting the handle does not guarantee that68        /// a read operation will return data. If using multiple69        /// reader threads, <see cref="NonBlockingRead"/> for70        /// performing non-blocking reads.71        /// </remarks>72        internal virtual WaitHandle ReadHandle73        {74            get75            {76#pragma warning disable 5650377                // disabled compiler warning as PSDataCollectionStream doesn't override this78                // and I didn't want code duplication.79                throw PSTraceSource.NewNotSupportedException();80#pragma warning restore 5650381            }82        }83 84        /// <summary>85        /// Waitable handle for callers to block until buffer space becomes available.86        /// </summary>87        /// <remarks>88        /// The handle is set when space becomes available for writing. For multiple89        /// writer threads writing to a bounded stream, the writer may still block90        /// if another thread fills the stream to capacity.91        /// </remarks>92        internal virtual WaitHandle WriteHandle93        {94            get95            {96#pragma warning disable 5650397                // disabled compiler warning as PSDataCollectionStream doesn't override this98                // and I didn't want code duplication.99                throw PSTraceSource.NewNotSupportedException();100#pragma warning restore 56503101            }102        }103 104        /// <summary>105        /// Determine if we are at the end of the stream.106        /// </summary>107        /// <remarks>108        /// EndOfPipeline is defined as the stream being closed and containing109        /// zero objects.  Readers check this to determine if any objects110        /// are in the stream.  Writers should check <see cref="IsOpen"/> to determine111        /// if the stream can be written to.112        /// </remarks>113        internal abstract bool EndOfPipeline { get; }114 115        /// <summary>116        /// Check if the stream is open for further writes.117        /// </summary>118        /// <returns>True if the stream is open, false if not.</returns>119        /// <remarks>120        /// IsOpen returns true until the first call to Close(). Writers should121        /// check IsOpen to determine if a write operation can be made.  Note that122        /// writers need to catch <see cref="PipelineClosedException"/>.123        /// <seealso cref="EndOfPipeline"/>124        /// </remarks>125        internal abstract bool IsOpen { get; }126 127        /// <summary>128        /// Returns the number of objects in the stream.129        /// </summary>130        internal abstract int Count { get; }131 132        /// <summary>133        /// Return a PipelineReader(object) for this stream.134        /// </summary>135        internal abstract PipelineReader<object> ObjectReader { get; }136 137        /// <summary>138        /// Return a PipelineReader(PSObject) for this stream.139        /// </summary>140        internal abstract PipelineReader<PSObject> PSObjectReader { get; }141 142        // 913921-2005/07/08 ObjectWriter can be retrieved on a closed stream143        /// <summary>144        /// Return an PipelineWriter for this stream.145        /// </summary>146        internal abstract PipelineWriter ObjectWriter { get; }147 148        #endregion149 150        #region Read Abstractions151 152        /// <summary>153        /// Read a single object from the stream.154        /// </summary>155        /// <returns>The next object in the stream or AutomationNull if EndOfPipeline is reached.</returns>156        /// <remarks>This method blocks if the stream is empty</remarks>157        internal virtual object Read()158        {159            throw PSTraceSource.NewNotSupportedException();160        }161 162        /// <summary>163        /// Read at most <paramref name="count"/> objects.164        /// </summary>165        /// <param name="count">The maximum number of objects to read.</param>166        /// <returns>The objects read.</returns>167        /// <exception cref="ArgumentOutOfRangeException">168        /// <paramref name="count"/> is less than 0169        /// </exception>170        /// <remarks>171        /// This method blocks if the number of objects in the stream is less than <paramref name="count"/>172        /// and the stream is not closed.173        ///174        /// If there are multiple reader threads, the objects returned175        /// to blocking reads Read(int count) and ReadToEnd()176        /// are not necessarily single blocks of objects added to the177        /// stream in that order.  For example, if ABCDEF are added to the178        /// stream, one reader may get ABDE and the other may get CF.179        /// Each reader reads items from the stream as they become available.180        /// Otherwise, if a maximum _capacity has been imposed, the writer181        /// and reader could become mutually deadlocked.182        ///183        /// When there are multiple blocked readers, any of the readers184        /// may get the next object(s) added.185        /// </remarks>186        internal virtual Collection<object> Read(int count)187        {188            throw PSTraceSource.NewNotSupportedException();189        }190 191        /// <summary>192        /// Blocks until the pipeline closes and reads all objects.193        /// </summary>194        /// <returns>A collection of zero or more objects.</returns>195        /// <remarks>196        /// If the stream is empty, a collection of size zero is returned.197        ///198        /// If there are multiple reader threads, the objects returned199        /// to blocking reads Read(int count) and ReadToEnd()200        /// are not necessarily single blocks of objects added to the201        /// stream in that order.  For example, if ABCDEF are added to the202        /// stream, one reader may get ABDE and the other may get CF.203        /// Each reader reads items from the stream as they become available.204        /// Otherwise, if a maximum _capacity has been imposed, the writer205        /// and reader could become mutually deadlocked.206        ///207        /// When there are multiple blocked readers, any of the readers208        /// may get the next object(s) added.209        /// </remarks>210        internal virtual Collection<object> ReadToEnd()211        {212            throw PSTraceSource.NewNotSupportedException();213        }214 215        /// <summary>216        /// Reads objects currently in the stream, but does not block.217        /// </summary>218        /// <returns>An array of zero or more objects.</returns>219        /// <remarks>220        /// This method performs a read of objects currently in the221        /// stream. The method will block until exclusive access to the222        /// stream is acquired.  If there are no objects in the stream,223        /// an empty array is returned.224        /// </remarks>225        /// <param name="maxRequested">226        /// Return no more than maxRequested objects.227        /// </param>228        /// <exception cref="ArgumentOutOfRangeException">229        /// <paramref name="maxRequested"/> is less than 0230        /// </exception>231        internal virtual Collection<object> NonBlockingRead(int maxRequested)232        {233            throw PSTraceSource.NewNotSupportedException();234        }235 236        /// <summary>237        /// Peek the next object.238        /// </summary>239        /// <returns>240        /// The next object in the stream or AutomationNull.Value if the stream is empty241        /// </returns>242        /// <exception cref="PipelineClosedException">The ObjectStream is closed.</exception>243        internal virtual object Peek()244        {245            throw PSTraceSource.NewNotSupportedException();246        }247 248        #endregion249 250        #region Write Abstractions251 252        /// <summary>253        /// Writes a object to the current position in the stream and254        /// advances the position within the stream by one object.255        /// </summary>256        /// <param name="value">The object to write to the stream.</param>257        /// <returns>258        /// One, if the write was successful, otherwise;259        /// zero if the stream was closed before the object could be written,260        /// or if the object was AutomationNull.Value.261        /// </returns>262        /// <exception cref="PipelineClosedException">263        /// The stream is closed264        /// </exception>265        /// <remarks>266        /// AutomationNull.Value is ignored267        /// </remarks>268        internal virtual int Write(object value)269        {270            return Write(value, false);271        }272 273        /// <summary>274        /// Write objects to the underlying stream.275        /// </summary>276        /// <param name="obj">Object or enumeration to read from.</param>277        /// <param name="enumerateCollection">278        /// If enumerateCollection is true, and <paramref name="obj"/>279        /// is an enumeration according to LanguagePrimitives.GetEnumerable,280        /// the objects in the enumeration will be unrolled and281        /// written separately.  Otherwise, <paramref name="obj"/>282        /// will be written as a single object.283        /// </param>284        /// <returns>The number of objects written.</returns>285        /// <exception cref="PipelineClosedException">286        /// The underlying stream is closed287        /// </exception>288        /// <remarks>289        /// If the enumeration contains elements equal to290        /// AutomationNull.Value, they are ignored.291        /// This can cause the return value to be less than the size of292        /// the collection.293        /// </remarks>294        internal virtual int Write(object obj, bool enumerateCollection)295        {296            throw PSTraceSource.NewNotSupportedException();297        }298 299        #endregion300 301        #region Close / Flush302 303        /// <summary>304        /// Close the stream.305        /// </summary>306        /// <remarks>307        /// Causes subsequent calls to IsOpen to return false and calls to308        /// a write operation to throw PipelineClosedException.309        /// All calls to Close() after the first call are silently ignored.310        /// </remarks>311        internal virtual void Close()312        {313            throw PSTraceSource.NewNotSupportedException();314        }315 316        /// <summary>317        /// Flush the data from the stream.  Closed streams may be flushed.318        /// </summary>319        internal virtual void Flush()320        {321            throw PSTraceSource.NewNotSupportedException();322        }323 324        #endregion325 326        #region IDisposable327 328        /// <summary>329        /// Public method for dispose.330        /// </summary>331        public void Dispose()332        {333            Dispose(true);334 335            GC.SuppressFinalize(this);336        }337 338        /// <summary>339        /// Release all resources.340        /// </summary>341        /// <param name="disposing">If true, release all managed resources.</param>342        protected abstract void Dispose(bool disposing);343 344        #endregion IDisposable345    }346 347    /// <summary>348    /// A FIFO memory based object stream.349    /// The purpose of this stream class is to provide the350    /// semantics of a unidirectional stream of objects351    /// between two threads using a dynamic memory buffer.352    /// </summary>353    /// <remarks>354    /// The stream may be bound or unbounded.  Bounded streams355    /// are created via passing a capacity to the constructor.356    /// Unbounded streams are created using the default constructor.357    ///358    /// The capacity of the stream can not be changed after359    /// construction.360    ///361    /// For bounded streams, attempts to write to the stream when362    /// the capacity has been reached causes the writer to block363    /// until objects are read.364    ///365    /// For unbounded streams, writers only block for the amount366    /// of time needed to acquire exclusive access to the367    /// stream.  Note that unbounded streams have a capacity of368    /// of Int32.MaxValue objects.  In theory, if this limit were369    /// reached, the stream would function as a bounded stream.370    ///371    /// This class is safe for multi-threaded use with the following372    /// side-effects:373    ///374    /// > For bounded streams, write operations are not guaranteed to375    /// be atomic.  If a write operation causes the capacity to be376    /// reached without writing all data, a partial write occurs and377    /// the writer blocks until data is read from the stream.378    ///379    /// > When multiple writer or reader threads are used, the order380    /// the reader or writer acquires a lock on the stream is381    /// undefined.  This means that the first call to write does not382    /// guarantee the writer will acquire a write lock first.  The first383    /// call to read does not guarantee the reader will acquire the384    /// read lock first.385    ///386    /// > Reads and writes may occur in any order. With a bounded387    /// stream, write operations between threads may also result in388    /// interleaved write operations.389    ///390    /// The result is that the order of data is only guaranteed if there is a391    /// single writer.392    /// </remarks>393    // 897230-2003/10/29-JonN marked sealed394    // 905990-2005/05/10-JonN Removed IDisposable395    internal sealed class ObjectStream : ObjectStreamBase, IDisposable396    {397        #region Private Fields398        /// <summary>399        /// Objects in the stream.400        /// </summary>401        // PERF-2003/08/22-JonN We should probably use Queue instead402        // PERF-2004/06/30-JonN Probably more efficient to use type403        //  Collection<object> as the underlying store404        private readonly List<object> _objects;405 406        /// <summary>407        /// Is the stream open or closed for writing?408        /// </summary>409        private bool _isOpen;410 411        #region Synchronization handles412        /// <summary>413        /// Read handle - signaled when data is ready to read.414        /// </summary>415        /// <remarks>416        /// This event may, on occasion, be signalled even when there is417        /// no data available.  If this happens, just wait again.418        /// Never wait on this event alone.  Since this is an AutoResetEvent,419        /// there is no way to definitely release all blocked threads when420        /// the stream is closed for reading.  Instead, use WaitAny on421        /// this handle and also _readClosedHandle.422        /// </remarks>423        private readonly AutoResetEvent _readHandle;424 425        /// <summary>426        /// Handle returned to callers for blocking on data ready.427        /// </summary>428        private ManualResetEvent _readWaitHandle;429 430        /// <summary>431        /// When this handle is set, the stream is closed for reading,432        /// so all blocked readers should be released.433        /// </summary>434        private readonly ManualResetEvent _readClosedHandle;435 436        /// <summary>437        /// Write handle - signaled with the number of objects in the438        /// stream becomes less than the maximum number of objects439        /// allowed in the stream.  <see cref="_capacity"/>440        /// </summary>441        /// <remarks>442        /// This event may, on occasion, be signalled even when there is443        /// no write buffer available.  If this happens, just wait again.444        /// Never wait on this event alone.  Since this is an AutoResetEvent,445        /// there is no way to definitely release all blocked threads when446        /// the stream is closed for writing.  Instead, use WaitAny on447        /// this handle and also _writeClosedHandle.448        /// </remarks>449        private readonly AutoResetEvent _writeHandle;450 451        /// <summary>452        /// Handle returned to callers for blocking until buffer space453        /// is available for write.454        /// </summary>455        private ManualResetEvent _writeWaitHandle;456 457        /// <summary>458        /// When this handle is set, the stream is closed for writing,459        /// so all blocked readers should be released.460        /// </summary>461        private readonly ManualResetEvent _writeClosedHandle;462        #endregion Synchronization handles463 464        /// <summary>465        /// The object reader for this stream.466        /// </summary>467        /// <remarks>468        /// This field is allocated on first demand and469        /// returned on subsequent calls.470        /// </remarks>471        private PipelineReader<object> _reader = null;472 473        /// <summary>474        /// The PSObject reader for this stream.475        /// </summary>476        /// <remarks>477        /// This field is allocated on first demand and478        /// returned on subsequent calls.479        /// </remarks>480        private PipelineReader<PSObject> _mshreader = null;481 482        /// <summary>483        /// The object writer for this stream.484        /// </summary>485        /// <remarks>486        /// This field is allocated on first demand and487        /// returned on subsequent calls.488        /// </remarks>489        private PipelineWriter _writer = null;490 491        /// <summary>492        /// Maximum number of objects allowed in the stream493        /// Note that this is not permitted to be more than Int32.MaxValue,494        /// since the underlying list has this limitation.495        /// </summary>496        private readonly int _capacity = Int32.MaxValue;497 498        /// <summary>499        /// This object is used to acquire an exclusive lock on the stream.500        /// </summary>501        /// <remarks>502        /// Note that we lock _monitorObject rather than "this" so that503        /// we are protected from outside code interfering in our504        /// critical section.  Thanks to Wintellect for the hint.505        /// </remarks>506        private readonly object _monitorObject = new object();507 508        /// <summary>509        /// Indicates if this stream has already been disposed.510        /// </summary>511        private bool _disposed = false;512 513        #endregion Private Fields514 515        #region Ctor516 517        /// <summary>518        /// Default constructor.519        /// </summary>520        /// <remarks>521        /// Constructs a stream with a maximum size of Int32.Max522        /// </remarks>523        internal ObjectStream()524            : this(Int32.MaxValue)525        {526        }527 528        /// <summary>529        /// Allocate the stream with an initial size.530        /// </summary>531        /// <param name="capacity">532        /// The maximum number of objects to allow in the buffer at a time.533        /// Note that this is not permitted to be more than Int32.MaxValue,534        /// since the underlying list has this limitation535        /// </param>536        /// <exception cref="ArgumentOutOfRangeException">537        /// <paramref name="_capacity"/> is less than or equal to zero538        /// <paramref name="_capacity"/> is greater than Int32.MaxValue539        /// </exception>540        internal ObjectStream(int capacity)541        {542            if (capacity <= 0 || capacity > Int32.MaxValue)543            {544                throw PSTraceSource.NewArgumentOutOfRangeException(nameof(capacity), capacity);545            }546 547            // the maximum number of objects to allow in the stream at a given time.548            _capacity = capacity;549 550            // event is not signaled since there is no data to read551            _readHandle = new AutoResetEvent(false);552 553            // event is signaled since there is available buffer space554            _writeHandle = new AutoResetEvent(true);555 556            // event is not signaled since the thread is still readable557            _readClosedHandle = new ManualResetEvent(false);558 559            // event is signaled since the thread is still writeable560            _writeClosedHandle = new ManualResetEvent(false);561 562            // the FIFO set of objects in the stream563            _objects = new List<object>();564 565            // Is the stream open?566            _isOpen = true;567        }568 569        #endregion Ctor570 571        #region internal properties572 573        /// <summary>574        /// Get the capacity of the stream.575        /// </summary>576        /// <value>577        /// The capacity of the stream.578        /// </value>579        /// <remarks>580        /// The capacity is the number of objects the stream may contain at one time.  Once this581        /// limit is reached, attempts to write into the stream block until buffer space582        /// becomes available.583        /// MaxCapacity cannot change, so we can skip the lock.584        /// </remarks>585        internal override int MaxCapacity586        {587            get588            {589                return _capacity;590            }591        }592 593        /// <summary>594        /// Waitable handle for callers to wait on until data ready to read.595        /// </summary>596        /// <remarks>597        /// The handle is set when data becomes available to read or598        /// when a partial read has completed.  If multiple readers599        /// are used, setting the handle does not guarantee that600        /// a read operation will return data. If using multiple601        /// reader threads, <see cref="NonBlockingRead"/> for602        /// performing non-blocking reads.603        /// </remarks>604        internal override WaitHandle ReadHandle605        {606            get607            {608                WaitHandle handle = null;609 610                lock (_monitorObject)611                {612                    // Create the handle signaled if there are objects in the stream613                    // or the stream has been closed.  The closed scenario addresses614                    // Pipeline readers that execute asynchronously.  Since the pipeline615                    // may complete with zero objects before the caller objects this616                    // handle, it will block indefinitely unless it is set.617                    _readWaitHandle ??= new ManualResetEvent(_objects.Count > 0 || !_isOpen);618 619                    handle = _readWaitHandle;620                }621 622                return handle;623            }624        }625 626        /// <summary>627        /// Waitable handle for callers to block until buffer space becomes available.628        /// </summary>629        /// <remarks>630        /// The handle is set when space becomes available for writing. For multiple631        /// writer threads writing to a bounded stream, the writer may still block632        /// if another thread fills the stream to capacity.633        /// </remarks>634        internal override WaitHandle WriteHandle635        {636            get637            {638                WaitHandle handle = null;639 640                lock (_monitorObject)641                {642                    _writeWaitHandle ??= new ManualResetEvent(_objects.Count < _capacity || !_isOpen);643 644                    handle = _writeWaitHandle;645                }646 647                return handle;648            }649        }650 651        /// <summary>652        /// Return a PipelineReader(object) for this stream.653        /// </summary>654        internal override PipelineReader<object> ObjectReader655        {656            get657            {658                PipelineReader<object> reader = null;659 660                lock (_monitorObject)661                {662                    // Always return an object reader, even if the stream663                    // is closed. This is to address requesting the object reader664                    // after calling Pipeline.Execute(). NOTE: If Execute completes665                    // without writing data to the output queue, the666                    // stream will be in the EndOfPipeline state because the667                    // stream is closed and has zero data.  Since this is a valid668                    // and expected execution path, we don't want to throw an exception.669                    _reader ??= new ObjectReader(this);670 671                    reader = _reader;672                }673 674                return reader;675            }676        }677 678        /// <summary>679        /// Return a PipelineReader(PSObject) for this stream.680        /// </summary>681        internal override PipelineReader<PSObject> PSObjectReader682        {683            get684            {685                PipelineReader<PSObject> reader = null;686 687                lock (_monitorObject)688                {689                    // Always return an object reader, even if the stream690                    // is closed. This is to address requesting the object reader691                    // after calling Pipeline.Execute(). NOTE: If Execute completes692                    // without writing data to the output queue, the693                    // stream will be in the EndOfPipeline state because the694                    // stream is closed and has zero data.  Since this is a valid695                    // and expected execution path, we don't want to throw an exception.696                    _mshreader ??= new PSObjectReader(this);697 698                    reader = _mshreader;699                }700 701                return reader;702            }703        }704 705        // 913921-2005/07/08 ObjectWriter can be retrieved on a closed stream706        /// <summary>707        /// Return an PipelineWriter for this stream.708        /// </summary>709        internal override PipelineWriter ObjectWriter710        {711            get712            {713                PipelineWriter writer = null;714 715                lock (_monitorObject)716                {717                    _writer ??= new ObjectWriter(this) as PipelineWriter;718 719                    writer = _writer;720                }721 722                return writer;723            }724        }725 726        /// <summary>727        /// Determine if we are at the end of the stream.728        /// </summary>729        /// <remarks>730        /// EndOfPipeline is defined as the stream being closed and containing731        /// zero objects.  Readers check this to determine if any objects732        /// are in the stream.  Writers should check <see cref="IsOpen"/> to determine733        /// if the stream can be written to.734        /// </remarks>735        internal override bool EndOfPipeline736        {737            get738            {739                bool endOfStream = true;740 741                lock (_monitorObject)742                {743                    endOfStream = (_objects.Count == 0 && !_isOpen);744                }745 746                return endOfStream;747            }748        }749 750        /// <summary>751        /// Check if the stream is open for further writes.752        /// </summary>753        /// <returns>True if the stream is open, false if not.</returns>754        /// <remarks>755        /// IsOpen returns true until the first call to Close(). Writers should756        /// check IsOpen to determine if a write operation can be made.  Note that757        /// writers need to catch <see cref="PipelineClosedException"/>.758        /// <seealso cref="EndOfPipeline"/>759        /// </remarks>760        internal override bool IsOpen761        {762            get763            {764                bool isOpen = true;765                // 2003/09/02-JonN Hitesh says that the access766                // of a bool variable is atomic so there is no need767                // for the lock.768                lock (_monitorObject)769                {770                    isOpen = _isOpen;771                }772 773                return isOpen;774            }775        }776 777        /// <summary>778        /// Returns the number of objects in the stream.779        /// </summary>780        internal override int Count781        {782            get783            {784                int count = 0;785 786                lock (_monitorObject)787                {788                    count = _objects.Count;789                }790 791                return count;792            }793        }794 795        #endregion internal properties796 797        #region private locking code798 799        /// <summary>800        /// Wait for data to be readable.801        /// </summary>802        /// <returns>True if EndOfPipeline is not reached.</returns>803        /// <remarks>804        /// WaitRead does not guarantee that data is present in the stream,805        /// only that data was added when the event was signaled.  Since there may be806        /// multiple readers, data may be removed from the stream807        /// before the caller has a chance to read the data.808        /// This method should never be called within a lock(_monitorObject).809        /// </remarks>810        private bool WaitRead()811        {812            if (!EndOfPipeline)813            {814                try815                {816                    WaitHandle[] ha = { _readHandle, _readClosedHandle };817                    WaitHandle.WaitAny(ha); // ignore return value818                }819                catch (ObjectDisposedException)820                {821                    // Since the _readHandle must be acquired outside822                    // a lock there's a chance that it was823                    // disposed after checking EndOfPipeline824                }825            }826 827            return !EndOfPipeline;828        }829 830        /// <summary>831        /// Wait for data to be writeable.832        /// </summary>833        /// <returns>True if the stream is writeable, otherwise; false.</returns>834        /// <remarks>835        /// WaitWrite does not guarantee that buffer space will be available in the stream836        /// when the caller attempts to write, only that buffer space was available837        /// when the event was signaled.838        /// This method should never be called within a lock(_monitorObject).839        /// </remarks>840        private bool WaitWrite()841        {842            if (IsOpen)843            {844                try845                {846                    WaitHandle[] ha = { _writeHandle, _writeClosedHandle };847                    WaitHandle.WaitAny(ha); // ignore return value848                }849                catch (ObjectDisposedException)850                {851                    // Since the _writeHandle must be acquired outside852                    // a lock there's a chance that it was853                    // disposed after checking IsOpen854                }855            }856 857            return IsOpen;858        }859 860        /// <summary>861        /// Utility method to signal handles and raise events862        /// in the consistent order.863        /// NOTE: Release the lock before raising events; otherwise,864        /// there is a possible deadlock during the readable event.865        /// </summary>866        /// <remarks>867        /// RaiseEvents is fairly idempotent, although it will signal868        /// DataReady every time.869        /// </remarks>870        private void RaiseEvents()871        {872            bool unblockReaders = true;873            bool unblockWriters = true;874            bool endOfStream = false;875            try876            {877                lock (_monitorObject)878                {879                    // External readers block only for open streams880                    // with no stored objects.  External writers block881                    // only for open streams with no free buffer space.882                    unblockReaders = (!_isOpen || (_objects.Count > 0));883                    unblockWriters = (!_isOpen || (_objects.Count < _capacity));884                    endOfStream = (!_isOpen && (_objects.Count == 0));885 886                    // I would prefer to set the ManualResetEvents outside887                    // of the lock, so that the unblocked thread would not888                    // immediately be re-blocked.  However, I am not889                    // confident that multiple sets/resets might not get890                    // out of order and leave the handle in the wrong state.891                    if (_readWaitHandle != null)892                    {893                        try894                        {895                            if (unblockReaders)896                            {897                                _readWaitHandle.Set();898                            }899                            else900                            {901                                _readWaitHandle.Reset();902                            }903                        }904                        catch (ObjectDisposedException)905                        {906                        }907                    }908 909                    if (_writeWaitHandle != null)910                    {911                        try912                        {913                            if (unblockWriters)914                            {915                                _writeWaitHandle.Set();916                            }917                            else918                            {919                                _writeWaitHandle.Reset();920                            }921                        }922                        catch (ObjectDisposedException)923                        {924                        }925                    }926                }927            }928            finally929            {930                // We prefer to set the AutoResetEvents outside of the lock,931                // so that the unblocked thread will not immediately be932                // re-blocked.  This works because setting the handle933                // is idempotent.934                if (unblockReaders)935                {936                    try937                    {938                        _readHandle.Set();939                    }940                    catch (ObjectDisposedException)941                    {942                    }943                }944 945                if (unblockWriters)946                {947                    try948                    {949                        _writeHandle.Set();950                    }951                    catch (ObjectDisposedException)952                    {953                    }954                }955 956                if (endOfStream)957                {958                    try959                    {960                        _readClosedHandle.Set();961                    }962                    catch (ObjectDisposedException)963                    {964                    }965                }966            }967 968            // This causes a synchronous call to the969            // client-provided handler.970            if (unblockReaders)971            {972                FireDataReadyEvent(this, EventArgs.Empty);973            }974#if (false)975            //976            // NOTE: Event delegates are called only after the internal977            // AutoResetEvents are set to ensure that an exception in an978            // event delegate does not leave the reset events in a bad979            // state.980            if (unblockWriters && WriteReady != null)981            {982                WriteReady (this, new EventArgs ());983            }984#endif985        }986 987        #endregion private locking code988 989        #region internal methods990 991        /// <summary>992        /// Flush the data from the stream.  Closed streams may be flushed.993        /// </summary>994        internal override void Flush()995        {996            bool raiseEvents = false;997 998            try999            {1000                lock (_monitorObject)1001                {1002                    if (_objects.Count > 0)1003                    {1004                        raiseEvents = true;1005                        _objects.Clear();1006                    }1007                }1008            }1009            finally1010            {1011                if (raiseEvents)1012                {1013                    RaiseEvents();1014                }1015            }1016        }1017 1018        /// <summary>1019        /// Close the stream.1020        /// </summary>1021        /// <remarks>1022        /// Causes subsequent calls to IsOpen to return false and calls to1023        /// a write operation to throw PipelineClosedException.1024        /// All calls to Close() after the first call are silently ignored.1025        /// </remarks>1026        internal override void Close()1027        {1028            bool raiseEvents = false;1029 1030            try1031            {1032                lock (_monitorObject)1033                {1034                    // if we transition from open to closed,1035                    // signal any blocking readers or writers1036                    // to ensure the close is seen.1037                    if (_isOpen)1038                    {1039                        raiseEvents = true;1040                        _isOpen = false;1041                    }1042                }1043            }1044            finally1045            {1046                if (raiseEvents)1047                {1048                    // RaiseEvents does not manage _writeClosedHandle1049                    try1050                    {1051                        _writeClosedHandle.Set();1052                    }1053                    catch (ObjectDisposedException)1054                    {1055                    }1056 1057                    RaiseEvents();1058                }1059            }1060        }1061 1062        #endregion internal methods1063 1064        #region Read Methods1065 1066        /// <summary>1067        /// Read a single object from the stream.1068        /// </summary>1069        /// <returns>The next object in the stream or AutomationNull if EndOfPipeline is reached.</returns>1070        /// <remarks>This method blocks if the stream is empty</remarks>1071        internal override object Read()1072        {1073            Collection<object> result = Read(1);1074            if (result.Count == 1)1075            {1076                return result[0];1077            }1078 1079            Diagnostics.Assert(result.Count == 0, "Invalid number of objects returned");1080 1081            return AutomationNull.Value;1082        }1083 1084        /// <summary>1085        /// Read at most <paramref name="count"/> objects.1086        /// </summary>1087        /// <param name="count">The maximum number of objects to read.</param>1088        /// <returns>The objects read.</returns>1089        /// <exception cref="ArgumentOutOfRangeException">1090        /// <paramref name="count"/> is less than 01091        /// </exception>1092        /// <remarks>1093        /// This method blocks if the number of objects in the stream is less than <paramref name="count"/>1094        /// and the stream is not closed.1095        ///1096        /// If there are multiple reader threads, the objects returned1097        /// to blocking reads Read(int count) and ReadToEnd()1098        /// are not necessarily single blocks of objects added to the1099        /// stream in that order.  For example, if ABCDEF are added to the1100        /// stream, one reader may get ABDE and the other may get CF.1101        /// Each reader reads items from the stream as they become available.1102        /// Otherwise, if a maximum _capacity has been imposed, the writer1103        /// and reader could become mutually deadlocked.1104        ///1105        /// When there are multiple blocked readers, any of the readers1106        /// may get the next object(s) added.1107        /// </remarks>1108        internal override Collection<object> Read(int count)1109        {1110            if (count < 0)1111            {1112                throw PSTraceSource.NewArgumentOutOfRangeException(nameof(count), count);1113            }1114 1115            if (count == 0)1116            {1117                return new Collection<object>();1118            }1119 1120            Collection<object> results = new Collection<object>();1121 1122            bool raiseEvents = false;1123            while ((count > 0) && WaitRead())1124            {1125                try1126                {1127                    lock (_monitorObject)1128                    {1129                        // double check to ensure data is ready1130                        if (_objects.Count == 0)1131                        {1132                            continue;    // wait some more1133                        }1134 1135                        raiseEvents = true;1136                        // NTRAID#Windows Out Of Band Releases-925566-2005/12/07-JonN1137                        int objectsAdded = 0;1138                        foreach (object o in _objects)1139                        {1140                            results.Add(o);1141                            objectsAdded++;1142                            if (--count <= 0)1143                                break;1144                        }1145 1146                        _objects.RemoveRange(0, objectsAdded);1147                    }1148                }1149                finally1150                {1151                    // Raise the appropriate read/write events outside the lock. This is1152                    // inside the while loop to ensure writers can have a chance to1153                    // write otherwise the reader will starve.1154                    // NOTE: This must occur in the finally block to ensure1155                    // the AutoResetEvents are left in the appropriate state, even1156                    // for error paths.1157                    if (raiseEvents)1158                    {1159                        RaiseEvents();1160                    }1161                }1162            }1163 1164            return results;1165        }1166 1167        /// <summary>1168        /// Blocks until the pipeline closes and reads all objects.1169        /// </summary>1170        /// <returns>A collection of zero or more objects.</returns>1171        /// <remarks>1172        /// If the stream is empty, a collection of size zero is returned.1173        ///1174        /// If there are multiple reader threads, the objects returned1175        /// to blocking reads Read(int count) and ReadToEnd()1176        /// are not necessarily single blocks of objects added to the1177        /// stream in that order.  For example, if ABCDEF are added to the1178        /// stream, one reader may get ABDE and the other may get CF.1179        /// Each reader reads items from the stream as they become available.1180        /// Otherwise, if a maximum _capacity has been imposed, the writer1181        /// and reader could become mutually deadlocked.1182        ///1183        /// When there are multiple blocked readers, any of the readers1184        /// may get the next object(s) added.1185        /// </remarks>1186        internal override Collection<object> ReadToEnd()1187        {1188            // NTRAID#Windows Out Of Band Releases-925566-2005/12/07-JonN1189            return Read(Int32.MaxValue);1190        }1191 1192        /// <summary>1193        /// Reads objects currently in the stream, but does not block.1194        /// </summary>1195        /// <returns>An array of zero or more objects.</returns>1196        /// <remarks>1197        /// This method performs a read of objects currently in the1198        /// stream. The method will block until exclusive access to the1199        /// stream is acquired.  If there are no objects in the stream,1200        /// an empty array is returned.

Showing the first 1,200 of 1944 lines. Download the file for the rest.