MegaBites-AI/Windows-powershell
0372
1// Copyright (c) Microsoft Corporation.2// Licensed under the MIT License.3 4using System.Collections.ObjectModel;5using System.Management.Automation.Runspaces;6using System.Runtime.InteropServices;7using System.Threading;8 9namespace System.Management.Automation.Internal10{11 /// <summary>12 /// A PipelineReader for an ObjectStream.13 /// </summary>14 /// <remarks>15 /// This class is not safe for multi-threaded operations.16 /// </remarks>17 internal abstract class ObjectReaderBase<T> : PipelineReader<T>, IDisposable18 {19 /// <summary>20 /// Construct with an existing ObjectStream.21 /// </summary>22 /// <param name="stream">The stream to read.</param>23 /// <exception cref="ArgumentNullException">Thrown if the specified stream is null.</exception>24 protected ObjectReaderBase([In, Out] ObjectStreamBase stream)25 {26 ArgumentNullException.ThrowIfNull(stream);27 28 _stream = stream;29 }30 31 #region Events32 33 /// <summary>34 /// Event fired when objects are added to the underlying stream.35 /// </summary>36 public override event EventHandler DataReady37 {38 add39 {40 lock (_monitorObject)41 {42 bool firstRegistrant = (InternalDataReady == null);43 InternalDataReady += value;44 if (firstRegistrant)45 {46 _stream.DataReady += this.OnDataReady;47 }48 }49 }50 51 remove52 {53 lock (_monitorObject)54 {55 InternalDataReady -= value;56 if (InternalDataReady == null)57 {58 _stream.DataReady -= this.OnDataReady;59 }60 }61 }62 }63 64 public event EventHandler InternalDataReady = null;65 66 #endregion Events67 68 #region Public Properties69 70 /// <summary>71 /// Waitable handle for caller's to block until data is ready to read from the underlying stream.72 /// </summary>73 public override WaitHandle WaitHandle74 {75 get76 {77 return _stream.ReadHandle;78 }79 }80 81 /// <summary>82 /// Check if the stream is closed and contains no data.83 /// </summary>84 /// <value>True if the stream is closed and contains no data, otherwise; false.</value>85 /// <remarks>86 /// Attempting to read from the underlying stream if EndOfPipeline is true returns87 /// zero objects.88 /// </remarks>89 public override bool EndOfPipeline90 {91 get92 {93 return _stream.EndOfPipeline;94 }95 }96 97 /// <summary>98 /// Check if the stream is open for further writes.99 /// </summary>100 /// <value>true if the underlying stream is open, otherwise; false.</value>101 /// <remarks>102 /// The underlying stream may be readable after it is closed if data remains in the103 /// internal buffer. Check <see cref="EndOfPipeline"/> to determine if104 /// the underlying stream is closed and contains no data.105 /// </remarks>106 public override bool IsOpen107 {108 get109 {110 return _stream.IsOpen;111 }112 }113 114 /// <summary>115 /// Returns the number of objects in the underlying stream.116 /// </summary>117 public override int Count118 {119 get120 {121 return _stream.Count;122 }123 }124 125 /// <summary>126 /// Get the capacity of the stream.127 /// </summary>128 /// <value>129 /// The capacity of the stream.130 /// </value>131 /// <remarks>132 /// The capacity is the number of objects that stream may contain at one time. Once this133 /// limit is reached, attempts to write into the stream block until buffer space134 /// becomes available.135 /// </remarks>136 public override int MaxCapacity137 {138 get139 {140 return _stream.MaxCapacity;141 }142 }143 144 #endregion Public Properties145 146 #region Public Methods147 148 /// <summary>149 /// Close the stream.150 /// </summary>151 /// <remarks>152 /// Causes subsequent calls to IsOpen to return false and calls to153 /// a write operation to throw an ObjectDisposedException.154 /// All calls to Close() after the first call are silently ignored.155 /// </remarks>156 /// <exception cref="ObjectDisposedException">157 /// The stream is already disposed158 /// </exception>159 public override void Close()160 {161 // 2003/09/02-JonN added call to close underlying stream162 _stream.Close();163 }164 165 #endregion Public Methods166 167 #region Private Methods168 169 /// <summary>170 /// Handle DataReady events from the underlying stream.171 /// </summary>172 /// <param name="sender">The stream raising the event.</param>173 /// <param name="args">Standard event args.</param>174 private void OnDataReady(object sender, EventArgs args)175 {176 // call any event handlers on this, replacing the177 // ObjectStream sender with 'this' since receivers178 // are expecting a PipelineReader<object>179 InternalDataReady.SafeInvoke(this, args);180 }181 182 #endregion Private Methods183 184 #region Private fields185 186 /// <summary>187 /// The underlying stream.188 /// </summary>189 /// <remarks>Can never be null</remarks>190 protected ObjectStreamBase _stream;191 192 /// <summary>193 /// This object is used to acquire an exclusive lock194 /// on event handler registration.195 /// </summary>196 /// <remarks>197 /// Note that we lock _monitorObject rather than "this" so that198 /// we are protected from outside code interfering in our199 /// critical section. Thanks to Wintellect for the hint.200 /// </remarks>201 private readonly object _monitorObject = new object();202 203 #endregion Private fields204 205 #region IDisposable206 207 /// <summary>208 /// Public method for dispose.209 /// </summary>210 public void Dispose()211 {212 Dispose(true);213 214 GC.SuppressFinalize(this);215 }216 217 /// <summary>218 /// Release all resources.219 /// </summary>220 /// <param name="disposing">If true, release all managed resources.</param>221 protected abstract void Dispose(bool disposing);222 223 #endregion IDisposable224 }225 226 /// <summary>227 /// A PipelineReader reading objects from an ObjectStream.228 /// </summary>229 /// <remarks>230 /// This class is not safe for multi-threaded operations.231 /// </remarks>232 internal class ObjectReader : ObjectReaderBase<object>233 {234 #region ctor235 /// <summary>236 /// Construct with an existing ObjectStream.237 /// </summary>238 /// <param name="stream">The stream to read.</param>239 /// <exception cref="ArgumentNullException">Thrown if the specified stream is null.</exception>240 public ObjectReader([In, Out] ObjectStream stream)241 : base(stream)242 { }243 #endregion ctor244 245 /// <summary>246 /// Read at most <paramref name="count"/> objects.247 /// </summary>248 /// <param name="count">The maximum number of objects to read.</param>249 /// <returns>The objects read.</returns>250 /// <remarks>251 /// This method blocks if the number of objects in the stream is less than <paramref name="count"/>252 /// and the stream is not closed.253 /// </remarks>254 public override Collection<object> Read(int count)255 {256 return _stream.Read(count);257 }258 259 /// <summary>260 /// Read a single object from the stream.261 /// </summary>262 /// <returns>The next object in the stream.</returns>263 /// <remarks>This method blocks if the stream is empty</remarks>264 public override object Read()265 {266 return _stream.Read();267 }268 269 /// <summary>270 /// Blocks until the pipeline closes and reads all objects.271 /// </summary>272 /// <returns>A collection of zero or more objects.</returns>273 /// <remarks>274 /// If the stream is empty, an empty collection is returned.275 /// </remarks>276 public override Collection<object> ReadToEnd()277 {278 return _stream.ReadToEnd();279 }280 281 /// <summary>282 /// Reads all objects currently in the stream, but does not block.283 /// </summary>284 /// <returns>A collection of zero or more objects.</returns>285 /// <remarks>286 /// This method performs a read of all objects currently in the287 /// stream. The method will block until exclusive access to the288 /// stream is acquired. If there are no objects in the stream,289 /// an empty collection is returned.290 /// </remarks>291 public override Collection<object> NonBlockingRead()292 {293 return _stream.NonBlockingRead(Int32.MaxValue);294 }295 296 /// <summary>297 /// Reads objects currently in the stream, but does not block.298 /// </summary>299 /// <returns>A collection of zero or more objects.</returns>300 /// <remarks>301 /// This method performs a read of objects currently in the302 /// stream. The method will block until exclusive access to the303 /// stream is acquired. If there are no objects in the stream,304 /// an empty collection is returned.305 /// </remarks>306 /// <param name="maxRequested">307 /// Return no more than maxRequested objects.308 /// </param>309 public override Collection<object> NonBlockingRead(int maxRequested)310 {311 return _stream.NonBlockingRead(maxRequested);312 }313 314 /// <summary>315 /// Peek the next object.316 /// </summary>317 /// <returns>The next object in the stream or ObjectStream.EmptyObject if the stream is empty.</returns>318 public override object Peek()319 {320 return _stream.Peek();321 }322 323 /// <summary>324 /// Release all resources.325 /// </summary>326 /// <param name="disposing">If true, release all managed resources.</param>327 protected override void Dispose(bool disposing)328 {329 if (disposing)330 {331 _stream.Close();332 }333 }334 }335 336 /// <summary>337 /// A PipelineReader reading PSObjects from an ObjectStream.338 /// </summary>339 /// <remarks>340 /// This class is not safe for multi-threaded operations.341 /// </remarks>342 internal class PSObjectReader : ObjectReaderBase<PSObject>343 {344 #region ctor345 /// <summary>346 /// Construct with an existing ObjectStream.347 /// </summary>348 /// <param name="stream">The stream to read.</param>349 /// <exception cref="ArgumentNullException">Thrown if the specified stream is null.</exception>350 public PSObjectReader([In, Out] ObjectStream stream)351 : base(stream)352 { }353 #endregion ctor354 355 /// <summary>356 /// Read at most <paramref name="count"/> objects.357 /// </summary>358 /// <param name="count">The maximum number of objects to read.</param>359 /// <returns>The objects read.</returns>360 /// <remarks>361 /// This method blocks if the number of objects in the stream is less than <paramref name="count"/>362 /// and the stream is not closed.363 /// </remarks>364 public override Collection<PSObject> Read(int count)365 {366 return MakePSObjectCollection(_stream.Read(count));367 }368 369 /// <summary>370 /// Read a single PSObject from the stream.371 /// </summary>372 /// <returns>The next PSObject in the stream.</returns>373 /// <remarks>This method blocks if the stream is empty</remarks>374 public override PSObject Read()375 {376 return MakePSObject(_stream.Read());377 }378 379 /// <summary>380 /// Blocks until the pipeline closes and reads all objects.381 /// </summary>382 /// <returns>A collection of zero or more objects.</returns>383 /// <remarks>384 /// If the stream is empty, an empty collection is returned.385 /// </remarks>386 public override Collection<PSObject> ReadToEnd()387 {388 return MakePSObjectCollection(_stream.ReadToEnd());389 }390 391 /// <summary>392 /// Reads all objects currently in the stream, but does not block.393 /// </summary>394 /// <returns>A collection of zero or more objects.</returns>395 /// <remarks>396 /// This method performs a read of all objects currently in the397 /// stream. The method will block until exclusive access to the398 /// stream is acquired. If there are no objects in the stream,399 /// an empty collection is returned.400 /// </remarks>401 public override Collection<PSObject> NonBlockingRead()402 {403 return MakePSObjectCollection(_stream.NonBlockingRead(Int32.MaxValue));404 }405 406 /// <summary>407 /// Reads objects currently in the stream, but does not block.408 /// </summary>409 /// <returns>A collection of zero or more objects.</returns>410 /// <remarks>411 /// This method performs a read of objects currently in the412 /// stream. The method will block until exclusive access to the413 /// stream is acquired. If there are no objects in the stream,414 /// an empty collection is returned.415 /// </remarks>416 /// <param name="maxRequested">417 /// Return no more than maxRequested objects.418 /// </param>419 public override Collection<PSObject> NonBlockingRead(int maxRequested)420 {421 return MakePSObjectCollection(_stream.NonBlockingRead(maxRequested));422 }423 424 /// <summary>425 /// Peek the next PSObject.426 /// </summary>427 /// <returns>The next PSObject in the stream or ObjectStream.EmptyObject if the stream is empty.</returns>428 public override PSObject Peek()429 {430 return MakePSObject(_stream.Peek());431 }432 433 /// <summary>434 /// Release all resources.435 /// </summary>436 /// <param name="disposing">If true, release all managed resources.</param>437 protected override void Dispose(bool disposing)438 {439 if (disposing)440 {441 _stream.Close();442 }443 }444 445 #region Private446 private static PSObject MakePSObject(object o)447 {448 if (o == null)449 return null;450 451 return PSObject.AsPSObject(o);452 }453 454 // It might ultimately be more efficient to455 // make ObjectStream generic and convert the objects to PSObject456 // before inserting them into the initial Collection, so that we457 // don't have to convert the collection later.458 private static Collection<PSObject> MakePSObjectCollection(459 Collection<object> coll)460 {461 if (coll == null)462 return null;463 Collection<PSObject> retval = new Collection<PSObject>();464 foreach (object o in coll)465 {466 retval.Add(MakePSObject(o));467 }468 469 return retval;470 }471 #endregion Private472 }473 474 /// <summary>475 /// A ObjectReader for a PSDataCollection ObjectStream.476 /// </summary>477 /// <remarks>478 /// PSDataCollection is introduced after 1.0. PSDataCollection is479 /// used to store data which can be used with different480 /// commands concurrently.481 /// Only Read() operation is supported currently.482 /// </remarks>483 internal class PSDataCollectionReader<T, TResult>484 : ObjectReaderBase<TResult>485 {486 #region Private Data487 488 private readonly PSDataCollectionEnumerator<T> _enumerator;489 490 #endregion491 492 #region ctor493 /// <summary>494 /// Construct with an existing ObjectStream.495 /// </summary>496 /// <param name="stream">The stream to read.</param>497 /// <exception cref="ArgumentNullException">Thrown if the specified stream is null.</exception>498 public PSDataCollectionReader(PSDataCollectionStream<T> stream)499 : base(stream)500 {501 System.Management.Automation.Diagnostics.Assert(stream.ObjectStore != null,502 "Stream should have a valid data store");503 _enumerator = (PSDataCollectionEnumerator<T>)stream.ObjectStore.GetEnumerator();504 }505 506 #endregion ctor507 508 /// <summary>509 /// This method is not supported.510 /// </summary>511 /// <param name="count">The maximum number of objects to read.</param>512 /// <returns>The objects read.</returns>513 public override Collection<TResult> Read(int count)514 {515 throw new NotSupportedException();516 }517 518 /// <summary>519 /// Read a single object from the stream.520 /// </summary>521 /// <returns>522 /// The next object in the buffer or AutomationNull if buffer is closed523 /// and data is not available.524 /// </returns>525 /// <remarks>526 /// This method blocks if the buffer is empty.527 /// </remarks>528 public override TResult Read()529 {530 object result = AutomationNull.Value;531 if (_enumerator.MoveNext())532 {533 result = _enumerator.Current;534 }535 536 return ConvertToReturnType(result);537 }538 539 /// <summary>540 /// This method is not supported.541 /// </summary>542 /// <returns></returns>543 /// <remarks></remarks>544 public override Collection<TResult> ReadToEnd()545 {546 throw new NotSupportedException();547 }548 549 /// <summary>550 /// This method is not supported.551 /// </summary>552 /// <returns></returns>553 /// <remarks></remarks>554 public override Collection<TResult> NonBlockingRead()555 {556 return NonBlockingRead(Int32.MaxValue);557 }558 559 /// <summary>560 /// This method is not supported.561 /// </summary>562 /// <returns></returns>563 /// <remarks></remarks>564 /// <param name="maxRequested">565 /// Return no more than maxRequested objects.566 /// </param>567 public override Collection<TResult> NonBlockingRead(int maxRequested)568 {569 if (maxRequested < 0)570 {571 throw PSTraceSource.NewArgumentOutOfRangeException(nameof(maxRequested), maxRequested);572 }573 574 if (maxRequested == 0)575 {576 return new Collection<TResult>();577 }578 579 Collection<TResult> result = new Collection<TResult>();580 int readCount = maxRequested;581 582 while (readCount > 0)583 {584 if (_enumerator.MoveNext(false))585 {586 result.Add(ConvertToReturnType(_enumerator.Current));587 continue;588 }589 590 break;591 }592 593 return result;594 }595 596 /// <summary>597 /// This method is not supported.598 /// </summary>599 /// <returns></returns>600 public override TResult Peek()601 {602 throw new NotSupportedException();603 }604 605 /// <summary>606 /// Release all resources.607 /// </summary>608 /// <param name="disposing">If true, release all managed resources.</param>609 protected override void Dispose(bool disposing)610 {611 if (disposing)612 {613 _stream.Close();614 }615 }616 617 private static TResult ConvertToReturnType(object inputObject)618 {619 Type resultType = typeof(TResult);620 if (typeof(PSObject) == resultType || typeof(object) == resultType)621 {622 TResult result;623 LanguagePrimitives.TryConvertTo(inputObject, out result);624 return result;625 }626 627 System.Management.Automation.Diagnostics.Assert(false,628 "ReturnType should be either object or PSObject only");629 throw PSTraceSource.NewNotSupportedException();630 }631 }632 633 /// <summary>634 /// A ObjectReader for a PSDataCollection ObjectStream.635 /// </summary>636 /// <remarks>637 /// PSDataCollection is introduced after 1.0. PSDataCollection is638 /// used to store data which can be used with different639 /// commands concurrently.640 /// Only Read() operation is supported currently.641 /// </remarks>642 internal class PSDataCollectionPipelineReader<T, TReturn>643 : ObjectReaderBase<TReturn>644 {645 #region Private Data646 647 private readonly PSDataCollection<T> _datastore;648 649 #endregion Private Data650 651 #region ctor652 /// <summary>653 /// Construct with an existing ObjectStream.654 /// </summary>655 /// <param name="stream">The stream to read.</param>656 /// <param name="computerName"></param>657 /// <param name="runspaceId"></param>658 internal PSDataCollectionPipelineReader(PSDataCollectionStream<T> stream,659 string computerName, Guid runspaceId)660 : base(stream)661 {662 System.Management.Automation.Diagnostics.Assert(stream.ObjectStore != null,663 "Stream should have a valid data store");664 _datastore = stream.ObjectStore;665 ComputerName = computerName;666 RunspaceId = runspaceId;667 }668 669 #endregion ctor670 671 /// <summary>672 /// Computer name passed in by the pipeline which673 /// created this reader.674 /// </summary>675 internal string ComputerName { get; }676 677 /// <summary>678 /// Runspace Id passed in by the pipeline which679 /// created this reader.680 /// </summary>681 internal Guid RunspaceId { get; }682 683 /// <summary>684 /// This method is not supported.685 /// </summary>686 /// <param name="count">The maximum number of objects to read.</param>687 /// <returns>The objects read.</returns>688 public override Collection<TReturn> Read(int count)689 {690 throw new NotSupportedException();691 }692 693 /// <summary>694 /// Read a single object from the stream.695 /// </summary>696 /// <returns>697 /// The next object in the buffer or AutomationNull if buffer is closed698 /// and data is not available.699 /// </returns>700 /// <remarks>701 /// This method blocks if the buffer is empty.702 /// </remarks>703 public override TReturn Read()704 {705 object result = AutomationNull.Value;706 if (_datastore.Count > 0)707 {708 Collection<T> resultCollection = _datastore.ReadAndRemove(1);709 710 // ReadAndRemove returns a Collection<T> type but we711 // just want the single object contained in the collection.712 if (resultCollection.Count == 1)713 {714 result = resultCollection[0];715 }716 }717 718 return ConvertToReturnType(result);719 }720 721 /// <summary>722 /// This method is not supported.723 /// </summary>724 /// <returns></returns>725 /// <remarks></remarks>726 public override Collection<TReturn> ReadToEnd()727 {728 throw new NotSupportedException();729 }730 731 /// <summary>732 /// This method is not supported.733 /// </summary>734 /// <returns></returns>735 /// <remarks></remarks>736 public override Collection<TReturn> NonBlockingRead()737 {738 return NonBlockingRead(Int32.MaxValue);739 }740 741 /// <summary>742 /// This method is not supported.743 /// </summary>744 /// <returns></returns>745 /// <remarks></remarks>746 /// <param name="maxRequested">747 /// Return no more than maxRequested objects.748 /// </param>749 public override Collection<TReturn> NonBlockingRead(int maxRequested)750 {751 if (maxRequested < 0)752 {753 throw PSTraceSource.NewArgumentOutOfRangeException(nameof(maxRequested), maxRequested);754 }755 756 if (maxRequested == 0)757 {758 return new Collection<TReturn>();759 }760 761 Collection<TReturn> results = new Collection<TReturn>();762 int readCount = maxRequested;763 764 while (readCount > 0)765 {766 if (_datastore.Count > 0)767 {768 results.Add(ConvertToReturnType((_datastore.ReadAndRemove(1))[0]));769 readCount--;770 continue;771 }772 773 break;774 }775 776 return results;777 }778 779 /// <summary>780 /// This method is not supported.781 /// </summary>782 /// <returns></returns>783 public override TReturn Peek()784 {785 throw new NotSupportedException();786 }787 788 /// <summary>789 /// Converts to the return type based on language primitives.790 /// </summary>791 /// <param name="inputObject">Input object to convert.</param>792 /// <returns>Input object converted to the specified return type.</returns>793 private static TReturn ConvertToReturnType(object inputObject)794 {795 Type resultType = typeof(TReturn);796 if (typeof(PSObject) == resultType || typeof(object) == resultType)797 {798 TReturn result;799 LanguagePrimitives.TryConvertTo(inputObject, out result);800 return result;801 }802 803 System.Management.Automation.Diagnostics.Assert(false,804 "ReturnType should be either object or PSObject only");805 throw PSTraceSource.NewNotSupportedException();806 }807 808 #region IDisposable809 810 /// <summary>811 /// Release all resources.812 /// </summary>813 /// <param name="disposing">If true, release all managed resources.</param>814 protected override void Dispose(bool disposing)815 {816 if (disposing)817 {818 _datastore.Dispose();819 }820 }821 822 #endregion IDisposable823 }824}825 