MegaBites-AI/Windows-powershell
0372
1// Copyright (c) Microsoft Corporation.2// Licensed under the MIT License.3 4using System.Collections.Generic;5using System.Collections.ObjectModel;6using System.IO;7using System.Management.Automation.Internal;8using System.Management.Automation.Tracing;9using System.Text;10using System.Xml;11 12using Dbg = System.Management.Automation.Diagnostics;13using TypeTable = System.Management.Automation.Runspaces.TypeTable;14 15namespace System.Management.Automation.Remoting16{17 /// <summary>18 /// This class is used to hold a fragment of remoting PSObject for transporting to remote computer.19 ///20 /// A large remoting PSObject will be broken into fragments. Each fragment has a ObjectId and a FragmentId.21 /// The first fragment has a StartFragment marker. The last fragment also an EndFragment marker.22 /// These fragments can be reassembled on the receiving23 /// end by sequencing the fragment ids.24 ///25 /// Currently control objects (Control-C for stopping a pipeline execution) is not26 /// really fragmented. These objects are small. They are just wrapped into a single27 /// fragment.28 /// </summary>29 internal class FragmentedRemoteObject30 {31 private byte[] _blob;32 private int _blobLength;33 34 /// <summary>35 /// SFlag stands for the IsStartFragment. It is the bit value in the binary encoding.36 /// </summary>37 internal const byte SFlag = 0x1;38 39 /// <summary>40 /// EFlag stands for the IsEndFragment. It is the bit value in the binary encoding.41 /// </summary>42 internal const byte EFlag = 0x2;43 44 /// <summary>45 /// HeaderLength is the total number of bytes in the binary encoding header.46 /// </summary>47 internal const int HeaderLength = 8 + 8 + 1 + 4;48 49 /// <summary>50 /// _objectIdOffset is the offset of the ObjectId in the binary encoding.51 /// </summary>52 private const int _objectIdOffset = 0;53 54 /// <summary>55 /// _fragmentIdOffset is the offset of the FragmentId in the binary encoding.56 /// </summary>57 private const int _fragmentIdOffset = 8;58 59 /// <summary>60 /// _flagsOffset is the offset of the byte in the binary encoding that contains the SFlag, EFlag and CFlag.61 /// </summary>62 private const int _flagsOffset = 16;63 64 /// <summary>65 /// _blobLengthOffset is the offset of the BlobLength in the binary encoding.66 /// </summary>67 private const int _blobLengthOffset = 17;68 69 /// <summary>70 /// _blobOffset is the offset of the Blob in the binary encoding.71 /// </summary>72 private const int _blobOffset = 21;73 74 #region Constructors75 76 /// <summary>77 /// Default Constructor.78 /// </summary>79 internal FragmentedRemoteObject()80 {81 }82 83 /// <summary>84 /// Used to construct a fragment of PSObject to be sent to remote computer.85 /// </summary>86 /// <param name="blob"></param>87 /// <param name="objectId">88 /// ObjectId of the fragment.89 /// Caller should make sure this is not less than 0.90 /// </param>91 /// <param name="fragmentId">92 /// FragmentId within the object.93 /// Caller should make sure this is not less than 0.94 /// </param>95 /// <param name="isEndFragment">96 /// true if this is a EndFragment.97 /// </param>98 internal FragmentedRemoteObject(byte[] blob, long objectId, long fragmentId,99 bool isEndFragment)100 {101 Dbg.Assert((blob != null) && (blob.Length != 0), "Cannot create a fragment for null or empty data.");102 Dbg.Assert(objectId >= 0, "Object Id cannot be < 0");103 Dbg.Assert(fragmentId >= 0, "Fragment Id cannot be < 0");104 105 ObjectId = objectId;106 FragmentId = fragmentId;107 108 IsStartFragment = fragmentId == 0;109 IsEndFragment = isEndFragment;110 111 _blob = blob;112 _blobLength = _blob.Length;113 }114 115 #endregion Constructors116 117 #region Data Fields being sent118 119 /// <summary>120 /// All fragments of the same PSObject have the same ObjectId.121 /// </summary>122 internal long ObjectId { get; set; }123 124 /// <summary>125 /// FragmentId starts from 0. It increases sequentially by an increment of 1.126 /// </summary>127 internal long FragmentId { get; set; }128 129 /// <summary>130 /// The first fragment of a PSObject.131 /// </summary>132 internal bool IsStartFragment { get; set; }133 134 /// <summary>135 /// The last fragment of a PSObject.136 /// </summary>137 internal bool IsEndFragment { get; set; }138 139 /// <summary>140 /// Blob length. This enables scenarios where entire byte[] is141 /// not filled for the fragment.142 /// </summary>143 internal int BlobLength144 {145 get146 {147 return _blobLength;148 }149 150 set151 {152 Dbg.Assert(value >= 0, "BlobLength cannot be less than 0.");153 _blobLength = value;154 }155 }156 157 /// <summary>158 /// This is the actual data in bytes form.159 /// </summary>160 internal byte[] Blob161 {162 get163 {164 return _blob;165 }166 167 set168 {169 Dbg.Assert(value != null, "Blob cannot be null");170 _blob = value;171 }172 }173 174 #endregion Data Fields being sent175 176 /// <summary>177 /// This method generate a binary encoding of the FragmentedRemoteObject as follows:178 /// ObjectId: 8 bytes as long, byte order is big-endian. this value can only be non-negative.179 /// FragmentId: 8 bytes as long, byte order is big-endian. this value can only be non-negative.180 /// FlagsByte: 1 byte:181 /// 0x1 if IsStartOfFragment is true: This is called S-flag.182 /// 0x2 if IsEndOfFragment is true: This is called the E-flag.183 /// 0x4 if IsControl is true: This is called the C-flag.184 ///185 /// The other bits are reserved for future use.186 /// Now they must be zero when sending,187 /// and they are ignored when receiving.188 /// BlobLength: 4 bytes as int, byte order is big-endian. this value can only be non-negative.189 /// Blob: BlobLength number of bytes.190 ///191 /// 0 1 2 3192 /// 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1193 /// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+194 /// | |195 /// +-+-+-+-+-+-+-+- ObjectId +-+-+-+-+-+-+-+-+196 /// | |197 /// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+198 /// | |199 /// +-+-+-+-+-+-+-+- FragmentId +-+-+-+-+-+-+-+-+200 /// | |201 /// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+202 /// |reserved |C|E|S|203 /// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+204 /// | BlobLength |205 /// +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+206 /// | Blob ...207 /// +-+-+-+-+-+-+-+-208 /// </summary>209 /// <returns>210 /// The binary encoded FragmentedRemoteObject to be ready to pass to WinRS Send API.211 /// </returns>212 internal byte[] GetBytes()213 {214 const int objectIdSize = 8; // number of bytes of long215 const int fragmentIdSize = 8; // number of bytes of long216 const int flagsSize = 1; // 1 byte for IsEndOfFrag and IsControl217 const int blobLengthSize = 4; // number of bytes of int218 219 int totalLength = objectIdSize + fragmentIdSize + flagsSize + blobLengthSize + BlobLength;220 221 byte[] result = new byte[totalLength];222 223 int idx = 0;224 225 // release build will optimize the calculation of the constants226 227 // ObjectId228 idx = _objectIdOffset;229 result[idx++] = (byte)((ObjectId >> (7 * 8)) & 0x7F); // sign bit is 0230 result[idx++] = (byte)((ObjectId >> (6 * 8)) & 0xFF);231 result[idx++] = (byte)((ObjectId >> (5 * 8)) & 0xFF);232 result[idx++] = (byte)((ObjectId >> (4 * 8)) & 0xFF);233 result[idx++] = (byte)((ObjectId >> (3 * 8)) & 0xFF);234 result[idx++] = (byte)((ObjectId >> (2 * 8)) & 0xFF);235 result[idx++] = (byte)((ObjectId >> 8) & 0xFF);236 result[idx++] = (byte)(ObjectId & 0xFF);237 238 // FragmentId239 idx = _fragmentIdOffset;240 result[idx++] = (byte)((FragmentId >> (7 * 8)) & 0x7F); // sign bit is 0241 result[idx++] = (byte)((FragmentId >> (6 * 8)) & 0xFF);242 result[idx++] = (byte)((FragmentId >> (5 * 8)) & 0xFF);243 result[idx++] = (byte)((FragmentId >> (4 * 8)) & 0xFF);244 result[idx++] = (byte)((FragmentId >> (3 * 8)) & 0xFF);245 result[idx++] = (byte)((FragmentId >> (2 * 8)) & 0xFF);246 result[idx++] = (byte)((FragmentId >> 8) & 0xFF);247 result[idx++] = (byte)(FragmentId & 0xFF);248 249 // E-flag and S-Flag250 idx = _flagsOffset;251 byte s_flag = IsStartFragment ? SFlag : (byte)0;252 byte e_flag = IsEndFragment ? EFlag : (byte)0;253 254 result[idx++] = (byte)(s_flag | e_flag);255 256 // BlobLength257 idx = _blobLengthOffset;258 result[idx++] = (byte)((BlobLength >> (3 * 8)) & 0xFF);259 result[idx++] = (byte)((BlobLength >> (2 * 8)) & 0xFF);260 result[idx++] = (byte)((BlobLength >> 8) & 0xFF);261 result[idx++] = (byte)(BlobLength & 0xFF);262 263 Array.Copy(_blob, 0, result, _blobOffset, BlobLength);264 265 return result;266 }267 268 /// <summary>269 /// Extract the objectId from a byte array, starting at the index indicated by270 /// startIndex parameter.271 /// </summary>272 /// <param name="fragmentBytes"></param>273 /// <param name="startIndex"></param>274 /// <returns>275 /// The objectId.276 /// </returns>277 /// <exception cref="ArgumentNullException">278 /// If fragmentBytes is null.279 /// </exception>280 /// <exception cref="ArgumentException">281 /// If startIndex is negative or fragmentBytes is not large enough to hold the entire header of282 /// a binary encoded FragmentedRemoteObject.283 /// </exception>284 internal static long GetObjectId(byte[] fragmentBytes, int startIndex)285 {286 Dbg.Assert(fragmentBytes != null, "fragmentBytes cannot be null");287 Dbg.Assert(fragmentBytes.Length >= HeaderLength, "not enough data to decode object id");288 long objectId = 0;289 290 int idx = startIndex + _objectIdOffset;291 292 objectId = (((long)fragmentBytes[idx++]) << (7 * 8)) & 0x7F00000000000000;293 objectId += (((long)fragmentBytes[idx++]) << (6 * 8)) & 0xFF000000000000;294 objectId += (((long)fragmentBytes[idx++]) << (5 * 8)) & 0xFF0000000000;295 objectId += (((long)fragmentBytes[idx++]) << (4 * 8)) & 0xFF00000000;296 objectId += (((long)fragmentBytes[idx++]) << (3 * 8)) & 0xFF000000;297 objectId += (((long)fragmentBytes[idx++]) << (2 * 8)) & 0xFF0000;298 objectId += (((long)fragmentBytes[idx++]) << 8) & 0xFF00;299 objectId += ((long)fragmentBytes[idx++]) & 0xFF;300 301 return objectId;302 }303 304 /// <summary>305 /// Extract the FragmentId from the byte array, starting at the index indicated by306 /// startIndex parameter.307 /// </summary>308 /// <param name="fragmentBytes"></param>309 /// <param name="startIndex"></param>310 /// <returns></returns>311 /// <exception cref="ArgumentNullException">312 /// If fragmentBytes is null.313 /// </exception>314 /// <exception cref="ArgumentException">315 /// If startIndex is negative or fragmentBytes is not large enough to hold the entire header of316 /// a binary encoded FragmentedRemoteObject.317 /// </exception>318 internal static long GetFragmentId(byte[] fragmentBytes, int startIndex)319 {320 Dbg.Assert(fragmentBytes != null, "fragmentBytes cannot be null");321 Dbg.Assert(fragmentBytes.Length >= HeaderLength, "not enough data to decode fragment id");322 long fragmentId = 0;323 int idx = startIndex + _fragmentIdOffset;324 325 fragmentId = (((long)fragmentBytes[idx++]) << (7 * 8)) & 0x7F00000000000000;326 fragmentId += (((long)fragmentBytes[idx++]) << (6 * 8)) & 0xFF000000000000;327 fragmentId += (((long)fragmentBytes[idx++]) << (5 * 8)) & 0xFF0000000000;328 fragmentId += (((long)fragmentBytes[idx++]) << (4 * 8)) & 0xFF00000000;329 fragmentId += (((long)fragmentBytes[idx++]) << (3 * 8)) & 0xFF000000;330 fragmentId += (((long)fragmentBytes[idx++]) << (2 * 8)) & 0xFF0000;331 fragmentId += (((long)fragmentBytes[idx++]) << 8) & 0xFF00;332 fragmentId += ((long)fragmentBytes[idx++]) & 0xFF;333 334 return fragmentId;335 }336 337 /// <summary>338 /// Extract the IsStartFragment value from the byte array, starting at the index indicated by339 /// startIndex parameter.340 /// </summary>341 /// <param name="fragmentBytes"></param>342 /// <param name="startIndex"></param>343 /// <returns>344 /// True is the S-flag is set in the encoding. Otherwise false.345 /// </returns>346 /// <exception cref="ArgumentNullException">347 /// If fragmentBytes is null.348 /// </exception>349 /// <exception cref="ArgumentException">350 /// If startIndex is negative or fragmentBytes is not large enough to hold the entire header of351 /// a binary encoded FragmentedRemoteObject.352 /// </exception>353 internal static bool GetIsStartFragment(byte[] fragmentBytes, int startIndex)354 {355 Dbg.Assert(fragmentBytes != null, "fragment cannot be null");356 Dbg.Assert(fragmentBytes.Length >= HeaderLength, "not enough data to decode if it is a start fragment.");357 358 if ((fragmentBytes[startIndex + _flagsOffset] & SFlag) != 0)359 {360 return true;361 }362 363 return false;364 }365 366 /// <summary>367 /// Extract the IsEndFragment value from the byte array, starting at the index indicated by368 /// startIndex parameter.369 /// </summary>370 /// <param name="fragmentBytes"></param>371 /// <param name="startIndex"></param>372 /// <returns>373 /// True if the E-flag is set in the encoding. Otherwise false.374 /// </returns>375 /// <exception cref="ArgumentNullException">376 /// If fragmentBytes is null.377 /// </exception>378 /// <exception cref="ArgumentException">379 /// If startIndex is negative or fragmentBytes is not large enough to hold the entire header of380 /// a binary encoded FragmentedRemoteObject.381 /// </exception>382 internal static bool GetIsEndFragment(byte[] fragmentBytes, int startIndex)383 {384 Dbg.Assert(fragmentBytes != null, "fragment cannot be null");385 Dbg.Assert(fragmentBytes.Length >= HeaderLength, "not enough data to decode if it is an end fragment.");386 387 if ((fragmentBytes[startIndex + _flagsOffset] & EFlag) != 0)388 {389 return true;390 }391 392 return false;393 }394 395 /// <summary>396 /// Extract the BlobLength value from the byte array, starting at the index indicated by397 /// startIndex parameter.398 /// </summary>399 /// <param name="fragmentBytes"></param>400 /// <param name="startIndex"></param>401 /// <returns>402 /// The BlobLength value.403 /// </returns>404 /// <exception cref="ArgumentNullException">405 /// If fragmentBytes is null.406 /// </exception>407 /// <exception cref="ArgumentException">408 /// If startIndex is negative or fragmentBytes is not large enough to hold the entire header of409 /// a binary encoded FragmentedRemoteObject.410 /// </exception>411 internal static int GetBlobLength(byte[] fragmentBytes, int startIndex)412 {413 Dbg.Assert(fragmentBytes != null, "fragment cannot be null");414 Dbg.Assert(fragmentBytes.Length >= HeaderLength, "not enough data to decode blob length.");415 416 int blobLength = 0;417 int idx = startIndex + _blobLengthOffset;418 419 blobLength += (((int)fragmentBytes[idx++]) << (3 * 8)) & 0x7F000000;420 blobLength += (((int)fragmentBytes[idx++]) << (2 * 8)) & 0xFF0000;421 blobLength += (((int)fragmentBytes[idx++]) << 8) & 0xFF00;422 blobLength += ((int)fragmentBytes[idx++]) & 0xFF;423 424 return blobLength;425 }426 }427 428 /// <summary>429 /// A stream used to store serialized data. This stream holds serialized data in the430 /// form of fragments. Every "fragment size" data will hold a blob identifying the fragment.431 /// The blob has "ObjectId","FragmentId","Properties like Start,End","BlobLength"432 /// </summary>433 internal class SerializedDataStream : Stream, IDisposable434 {435 [TraceSource("SerializedDataStream", "SerializedDataStream")]436 private static readonly PSTraceSource s_trace = PSTraceSource.GetTracer("SerializedDataStream", "SerializedDataStream");437 #region Global Constants438 439 private static long s_objectIdSequenceNumber = 0;440 441 #endregion442 443 #region Private Data444 445 private bool _isEntered;446 private readonly FragmentedRemoteObject _currentFragment;447 private long _fragmentId;448 449 private readonly int _fragmentSize;450 private readonly object _syncObject;451 private bool _isDisposed;452 private readonly bool _notifyOnWriteFragmentImmediately;453 454 // MemoryStream does not dynamically resize as data is read. This will waste455 // lot of memory as data sent on the network will still be there in memory.456 // To avoid this a queue of memory streams (each stream is of fragmentsize)457 // is created..so after data is sent the MemoryStream is disposed there by458 // clearing resources.459 private readonly Queue<MemoryStream> _queuedStreams;460 private MemoryStream _writeStream;461 private MemoryStream _readStream;462 private int _writeOffset;463 private int _readOffSet;464 private long _length;465 466 /// <summary>467 /// Callback that is called once a fragmented data is available.468 /// </summary>469 /// <param name="data">470 /// Data that resulted in this callback.471 /// </param>472 /// <param name="isEndFragment">473 /// true if data represents EndFragment of an object.474 /// </param>475 internal delegate void OnDataAvailableCallback(byte[] data, bool isEndFragment);476 477 private OnDataAvailableCallback _onDataAvailableCallback;478 479 #endregion480 481 #region Constructor482 483 /// <summary>484 /// Creates a stream to hold serialized data.485 /// </summary>486 /// <param name="fragmentSize">487 /// fragmentSize to be used while creating fragment boundaries.488 /// </param>489 internal SerializedDataStream(int fragmentSize)490 {491 s_trace.WriteLine("Creating SerializedDataStream with fragmentsize : {0}", fragmentSize);492 Dbg.Assert(fragmentSize > 0, "fragmentsize should be greater than 0.");493 _syncObject = new object();494 _currentFragment = new FragmentedRemoteObject();495 _queuedStreams = new Queue<MemoryStream>();496 _fragmentSize = fragmentSize;497 }498 499 /// <summary>500 /// Use this constructor carefully. This will not write data into internal501 /// streams. Instead this will make the SerializedDataStream call the502 /// callback whenever a fragmented data is available. It is upto the caller503 /// to figure out what to do with the data.504 /// </summary>505 /// <param name="fragmentSize">506 /// fragmentSize to be used while creating fragment boundaries.507 /// </param>508 /// <param name="callbackToNotify">509 /// If this is not null, then callback will get notified whenever fragmented510 /// data is available. Read() will return null in this case always.511 /// </param>512 internal SerializedDataStream(int fragmentSize,513 OnDataAvailableCallback callbackToNotify) : this(fragmentSize)514 {515 if (callbackToNotify != null)516 {517 _notifyOnWriteFragmentImmediately = true;518 _onDataAvailableCallback = callbackToNotify;519 }520 }521 522 #endregion523 524 #region Internal methods / Protected overrides525 526 /// <summary>527 /// Start using the stream exclusively (to write data). The stream can be entered only once.528 /// If you want to Enter again, first Exit and then Enter.529 /// This method is not thread-safe.530 /// </summary>531 internal void Enter()532 {533 Dbg.Assert(!_isEntered, "Stream is already entered. You cannot enter into stream again.");534 _isEntered = true;535 _fragmentId = 0;536 537 // Initialize the current fragment538 _currentFragment.ObjectId = GetObjectId();539 _currentFragment.FragmentId = _fragmentId;540 _currentFragment.IsStartFragment = true;541 _currentFragment.BlobLength = 0;542 _currentFragment.Blob = new byte[_fragmentSize];543 }544 545 /// <summary>546 /// Notify that the stream is not used to write anymore.547 /// This method is not thread-safe.548 /// </summary>549 internal void Exit()550 {551 _isEntered = false;552 // write left over data553 if (_currentFragment.BlobLength > 0)554 {555 // this is endfragment...as we are in Exit556 _currentFragment.IsEndFragment = true;557 WriteCurrentFragmentAndReset();558 }559 }560 561 /// <summary>562 /// Writes a block of bytes to the current stream using data read from buffer.563 /// The base MemoryStream is written to only if "FragmentSize" is reached.564 /// </summary>565 /// <param name="buffer">566 /// The buffer to read data from.567 /// </param>568 /// <param name="offset">569 /// The byte offset in buffer at which to begin writing from.570 /// </param>571 /// <param name="count">572 /// The maximum number of bytes to write.573 /// </param>574 public override void Write(byte[] buffer, int offset, int count)575 {576 Dbg.Assert(_isEntered, "Stream should be Entered before writing into.");577 578 int offsetToReadFrom = offset;579 int amountLeft = count;580 581 while (amountLeft > 0)582 {583 int dataLeftInTheFragment = _fragmentSize - FragmentedRemoteObject.HeaderLength - _currentFragment.BlobLength;584 if (dataLeftInTheFragment > 0)585 {586 int amountToWriteIntoFragment = (amountLeft > dataLeftInTheFragment) ? dataLeftInTheFragment : amountLeft;587 amountLeft -= amountToWriteIntoFragment;588 589 // Write data into fragment590 Array.Copy(buffer, offsetToReadFrom, _currentFragment.Blob, _currentFragment.BlobLength, amountToWriteIntoFragment);591 _currentFragment.BlobLength += amountToWriteIntoFragment;592 offsetToReadFrom += amountToWriteIntoFragment;593 594 // write only if amountLeft is more than 0. I dont write if amountLeft is 0 as we are not595 // sure if the fragment is EndFragment..we will know this only in Exit.596 if (amountLeft > 0)597 {598 WriteCurrentFragmentAndReset();599 }600 }601 else602 {603 WriteCurrentFragmentAndReset();604 }605 }606 }607 608 /// <summary>609 /// Writes a byte to the current stream.610 /// </summary>611 /// <param name="value"></param>612 public override void WriteByte(byte value)613 {614 Dbg.Assert(_isEntered, "Stream should be Entered before writing into.");615 byte[] buffer = new byte[1];616 buffer[0] = value;617 Write(buffer, 0, 1);618 }619 620 /// <summary>621 /// Returns a byte[] which holds data of fragment size (or) serialized data of622 /// one object, which ever is greater. If data is not currently available, then623 /// the callback is registered and called whenever the data is available.624 /// </summary>625 /// <param name="callback">626 /// callback to call once the data becomes available.627 /// </param>628 /// <returns>629 /// a byte[] holding data read from the stream630 /// </returns>631 internal byte[] ReadOrRegisterCallback(OnDataAvailableCallback callback)632 {633 lock (_syncObject)634 {635 if (_length <= 0)636 {637 _onDataAvailableCallback = callback;638 return null;639 }640 641 int bytesToRead = _length > _fragmentSize ? _fragmentSize : (int)_length;642 byte[] result = new byte[bytesToRead];643 Read(result, 0, bytesToRead);644 return result;645 }646 }647 648 /// <summary>649 /// Read the currently accumulated data in queued memory streams.650 /// </summary>651 /// <returns></returns>652 internal byte[] Read()653 {654 lock (_syncObject)655 {656 if (_isDisposed)657 {658 return null;659 }660 661 int bytesToRead = _length > _fragmentSize ? _fragmentSize : (int)_length;662 if (bytesToRead > 0)663 {664 byte[] result = new byte[bytesToRead];665 Read(result, 0, bytesToRead);666 return result;667 }668 else669 {670 return null;671 }672 }673 }674 675 /// <summary>676 /// </summary>677 /// <param name="buffer"></param>678 /// <param name="offset"></param>679 /// <param name="count"></param>680 /// <returns></returns>681 public override int Read(byte[] buffer, int offset, int count)682 {683 int offSetToWriteTo = offset;684 int dataWritten = 0;685 Collection<MemoryStream> memoryStreamsToDispose = new Collection<MemoryStream>();686 MemoryStream prevReadStream = null;687 688 lock (_syncObject)689 {690 // technically this should throw an exception..but remoting callstack691 // is optimized ie., we are not locking in every layer (in powershell)692 // to save on performance..as a result there may be cases where693 // upper layer is trying to add stuff and stream is disposed while694 // adding stuff.695 if (_isDisposed)696 {697 return 0;698 }699 700 while (dataWritten < count)701 {702 if (_readStream == null)703 {704 if (_queuedStreams.Count > 0)705 {706 _readStream = _queuedStreams.Dequeue();707 if ((!_readStream.CanRead) || (prevReadStream == _readStream))708 {709 // if the stream is disposed CanRead returns false710 // this will happen if a Write enqueues the stream711 // and a Read reads the data without dequeuing712 _readStream = null;713 continue;714 }715 }716 else717 {718 _readStream = _writeStream;719 }720 721 Dbg.Assert(_readStream.Length > 0, "Not enough data to read.");722 _readOffSet = 0;723 }724 725 _readStream.Position = _readOffSet;726 int result = _readStream.Read(buffer, offSetToWriteTo, count - dataWritten);727 s_trace.WriteLine("Read {0} data from readstream: {1}", result, _readStream.GetHashCode());728 dataWritten += result;729 offSetToWriteTo += result;730 _readOffSet += result;731 _length -= result;732 733 // dispose only if we dont read from the current write stream.734 if ((_readStream.Capacity == _readOffSet) && (_readStream != _writeStream))735 {736 s_trace.WriteLine("Adding readstream {0} to dispose collection.", _readStream.GetHashCode());737 memoryStreamsToDispose.Add(_readStream);738 prevReadStream = _readStream;739 _readStream = null;740 }741 }742 }743 744 // Dispose the memory streams outside of the lock745 foreach (MemoryStream streamToDispose in memoryStreamsToDispose)746 {747 s_trace.WriteLine("Disposing stream: {0}", streamToDispose.GetHashCode());748 streamToDispose.Dispose();749 }750 751 return dataWritten;752 }753 754 private void WriteCurrentFragmentAndReset()755 {756 // log trace of the fragment757 PSEtwLog.LogAnalyticVerbose(758 PSEventId.SentRemotingFragment, PSOpcode.Send, PSTask.None,759 PSKeyword.Transport | PSKeyword.UseAlwaysAnalytic,760 (Int64)(_currentFragment.ObjectId),761 (Int64)(_currentFragment.FragmentId),762 _currentFragment.IsStartFragment ? 1 : 0,763 _currentFragment.IsEndFragment ? 1 : 0,764 (UInt32)(_currentFragment.BlobLength),765 new PSETWBinaryBlob(_currentFragment.Blob, 0, _currentFragment.BlobLength));766 767 // finally write into memory stream768 byte[] data = _currentFragment.GetBytes();769 int amountLeft = data.Length;770 int offSetToReadFrom = 0;771 772 // user asked us to notify immediately..so no need773 // to write into memory stream..instead give the774 // data directly to user and let him figure out what to do.775 // This will save write + read + dispose!!776 if (!_notifyOnWriteFragmentImmediately)777 {778 lock (_syncObject)779 {780 // technically this should throw an exception..but remoting callstack781 // is optimized ie., we are not locking in every layer (in powershell)782 // to save on performance..as a result there may be cases where783 // upper layer is trying to add stuff and stream is disposed while784 // adding stuff.785 if (_isDisposed)786 {787 return;788 }789 790 if (_writeStream == null)791 {792 _writeStream = new MemoryStream(_fragmentSize);793 s_trace.WriteLine("Created write stream: {0}", _writeStream.GetHashCode());794 _writeOffset = 0;795 }796 797 while (amountLeft > 0)798 {799 int dataLeftInWriteStream = _writeStream.Capacity - _writeOffset;800 if (dataLeftInWriteStream == 0)801 {802 // enqueue the current write stream and create a new one.803 EnqueueWriteStream();804 dataLeftInWriteStream = _writeStream.Capacity - _writeOffset;805 }806 807 int amountToWriteIntoStream = (amountLeft > dataLeftInWriteStream) ? dataLeftInWriteStream : amountLeft;808 amountLeft -= amountToWriteIntoStream;809 // write data810 _writeStream.Position = _writeOffset;811 _writeStream.Write(data, offSetToReadFrom, amountToWriteIntoStream);812 offSetToReadFrom += amountToWriteIntoStream;813 _writeOffset += amountToWriteIntoStream;814 _length += amountToWriteIntoStream;815 }816 }817 }818 819 // call the callback since we have data available820 _onDataAvailableCallback?.Invoke(data, _currentFragment.IsEndFragment);821 822 // prepare a new fragment823 _currentFragment.FragmentId = ++_fragmentId;824 _currentFragment.IsStartFragment = false;825 _currentFragment.IsEndFragment = false;826 _currentFragment.BlobLength = 0;827 _currentFragment.Blob = new byte[_fragmentSize];828 }829 830 private void EnqueueWriteStream()831 {832 s_trace.WriteLine("Queuing write stream: {0} Length: {1} Capacity: {2}",833 _writeStream.GetHashCode(), _writeStream.Length, _writeStream.Capacity);834 _queuedStreams.Enqueue(_writeStream);835 836 _writeStream = new MemoryStream(_fragmentSize);837 _writeOffset = 0;838 s_trace.WriteLine("Created write stream: {0}", _writeStream.GetHashCode());839 }840 /// <summary>841 /// This method provides a thread safe way to get an object id.842 /// </summary>843 /// <returns>844 /// An object Id in integer.845 /// </returns>846 private static long GetObjectId()847 {848 return System.Threading.Interlocked.Increment(ref s_objectIdSequenceNumber);849 }850 851 #endregion852 853 #region Disposable Overrides854 855 protected override void Dispose(bool disposing)856 {857 if (disposing)858 {859 lock (_syncObject)860 {861 foreach (MemoryStream streamToDispose in _queuedStreams)862 {863 // make sure we dispose only once.864 if (streamToDispose.CanRead)865 {866 streamToDispose.Dispose();867 }868 }869 870 if ((_readStream != null) && (_readStream.CanRead))871 {872 _readStream.Dispose();873 }874 875 if ((_writeStream != null) && (_writeStream.CanRead))876 {877 _writeStream.Dispose();878 }879 880 _isDisposed = true;881 }882 }883 }884 885 #endregion886 887 #region Stream Overrides888 889 /// <summary>890 /// </summary>891 public override bool CanRead { get { return true; } }892 893 /// <summary>894 /// </summary>895 public override bool CanSeek { get { return false; } }896 897 /// <summary>898 /// </summary>899 public override bool CanWrite { get { return true; } }900 /// <summary>901 /// Gets the length of the stream in bytes.902 /// </summary>903 public override long Length { get { return _length; } }904 905 /// <summary>906 /// </summary>907 public override long Position908 {909 get { throw new NotSupportedException(); }910 911 set { throw new NotSupportedException(); }912 }913 /// <summary>914 /// This is a No-Op intentionally as there is nothing915 /// to flush.916 /// </summary>917 public override void Flush()918 {919 }920 921 /// <summary>922 /// </summary>923 /// <param name="offset"></param>924 /// <param name="origin"></param>925 /// <returns></returns>926 public override long Seek(long offset, SeekOrigin origin)927 {928 throw new NotSupportedException();929 }930 931 /// <summary>932 /// </summary>933 /// <param name="value"></param>934 public override void SetLength(long value)935 {936 throw new NotSupportedException();937 }938 939 #endregion940 941 #region IDisposable Members942 943 private bool _disposed = false;944 945 public new void Dispose()946 {947 if (!_disposed)948 {949 GC.SuppressFinalize(this);950 _disposed = true;951 }952 953 base.Dispose();954 }955 956 #endregion957 }958 959 /// <summary>960 /// This class performs the fragmentation as well as defragmentation operations of large objects to be sent961 /// to the other side. A large remoting PSObject will be broken into fragments. Each fragment has a ObjectId962 /// and a FragmentId. The last fragment also has an end of fragment marker. These fragments can be reassembled963 /// on the receiving end by sequencing the fragment ids.964 /// </summary>965 internal class Fragmentor966 {967 #region Global Constants968 private static readonly UTF8Encoding s_utf8Encoding = new UTF8Encoding();969 // This const defines the default depth to be used for serializing objects for remoting.970 private const int SerializationDepthForRemoting = 1;971 #endregion972 973 private int _fragmentSize;974 private readonly SerializationContext _serializationContext;975 976 #region Constructor977 978 /// <summary>979 /// Constructor which initializes fragmentor with FragmentSize.980 /// </summary>981 /// <param name="fragmentSize">982 /// size of each fragment983 /// </param>984 /// <param name="cryptoHelper"></param>985 internal Fragmentor(int fragmentSize, PSRemotingCryptoHelper cryptoHelper)986 {987 Dbg.Assert(fragmentSize > 0, "fragment size cannot be less than 0.");988 _fragmentSize = fragmentSize;989 _serializationContext = new SerializationContext(990 SerializationDepthForRemoting,991 SerializationOptions.RemotingOptions,992 cryptoHelper);993 DeserializationContext = new DeserializationContext(994 DeserializationOptions.RemotingOptions,995 cryptoHelper);996 }997 998 #endregion999 1000 /// <summary>1001 /// The method performs the fragmentation operation.1002 /// All fragments of the same object have the same ObjectId.1003 /// All fragments of the same object have the same ObjectId.1004 /// Each fragment has its own Fragment Id. Fragment Id always starts from zero (0),1005 /// and increments sequentially with an increment of 1.1006 /// The last fragment is indicated by an End of Fragment marker.1007 /// </summary>1008 /// <param name="obj">1009 /// The object to be fragmented. Caller should make sure this is not null.1010 /// </param>1011 /// <param name="dataToBeSent">1012 /// Caller specified dataToStore to which the fragments are added1013 /// one-by-one1014 /// </param>1015 internal void Fragment<T>(RemoteDataObject<T> obj, SerializedDataStream dataToBeSent)1016 {1017 Dbg.Assert(obj != null, "Cannot fragment a null object");1018 Dbg.Assert(dataToBeSent != null, "SendDataCollection cannot be null");1019 1020 dataToBeSent.Enter();1021 try1022 {1023 obj.Serialize(dataToBeSent, this);1024 }1025 finally1026 {1027 dataToBeSent.Exit();1028 }1029 }1030 1031 /// <summary>1032 /// The deserialization context used by this fragmentor. DeserializationContext1033 /// controls the amount of memory a deserializer can use and other things.1034 /// </summary>1035 internal DeserializationContext DeserializationContext { get; }1036 1037 /// <summary>1038 /// The size limit of the fragmented object.1039 /// </summary>1040 internal int FragmentSize1041 {1042 get1043 {1044 return _fragmentSize;1045 }1046 1047 set1048 {1049 Dbg.Assert(value > 0, "FragmentSize cannot be less than 0.");1050 _fragmentSize = value;1051 }1052 }1053 1054 /// <summary>1055 /// TypeTable used for Serialization/Deserialization.1056 /// </summary>1057 internal TypeTable TypeTable { get; set; }1058 1059 /// <summary>1060 /// Serialize an PSObject into a byte array.1061 /// </summary>1062 internal void SerializeToBytes(object obj, Stream streamToWriteTo)1063 {1064 Dbg.Assert(obj != null, "Cannot serialize a null object");1065 Dbg.Assert(streamToWriteTo != null, "Stream to write to cannot be null");1066 1067 XmlWriterSettings xmlSettings = new XmlWriterSettings();1068 xmlSettings.CheckCharacters = false;1069 xmlSettings.Indent = false;1070 // we dont want the underlying stream to be closed as we expect1071 // the stream to be usable after this call.1072 xmlSettings.CloseOutput = false;1073 xmlSettings.Encoding = UTF8Encoding.UTF8;1074 xmlSettings.NewLineHandling = NewLineHandling.None;1075 1076 xmlSettings.OmitXmlDeclaration = true;1077 xmlSettings.ConformanceLevel = ConformanceLevel.Fragment;1078 1079 using (XmlWriter xmlWriter = XmlWriter.Create(streamToWriteTo, xmlSettings))1080 {1081 Serializer serializer = new Serializer(xmlWriter, _serializationContext);1082 serializer.TypeTable = TypeTable;1083 serializer.Serialize(obj);1084 serializer.Done();1085 xmlWriter.Flush();1086 }1087 1088 return;1089 }1090 1091 /// <summary>1092 /// Converts the bytes back to PSObject.1093 /// </summary>1094 /// <param name="serializedDataStream">1095 /// The bytes to be deserialized.1096 /// </param>1097 /// <returns>1098 /// The deserialized object.1099 /// </returns>1100 /// <exception cref="PSRemotingDataStructureException">1101 /// If the deserialized object is null.1102 /// </exception>1103 internal PSObject DeserializeToPSObject(Stream serializedDataStream)1104 {1105 Dbg.Assert(serializedDataStream != null, "Cannot Deserialize null data");1106 Dbg.Assert(serializedDataStream.Length != 0, "Cannot Deserialize empty data");1107 1108 object result = null;1109 using (XmlReader xmlReader = XmlReader.Create(serializedDataStream, InternalDeserializer.XmlReaderSettingsForCliXml))1110 {1111 Deserializer deserializer = new Deserializer(xmlReader, DeserializationContext);1112 deserializer.TypeTable = TypeTable;1113 result = deserializer.Deserialize();1114 deserializer.Done();1115 }1116 1117 if (result == null)1118 {1119 // cannot be null.1120 throw new PSRemotingDataStructureException(RemotingErrorIdStrings.DeserializedObjectIsNull);1121 }1122 1123 return PSObject.AsPSObject(result);1124 }1125 }1126}1127 