MegaBites-AI/Windows-powershell
0372
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.