Team Ai
Datasetpublic

MegaBites-AI/Windows-powershell

sourceHugging Facemitupdated 6mo agoView on Hugging Face
0likes372downloads
fragmentor.cs1127 linesDownload Raw Back to common
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