Team Ai
Datasetpublic

MegaBites-AI/Windows-powershell

sourceHugging Facemitupdated 6mo agoView on Hugging Face
0likes372downloads
ThrottlingJob.cs1275 linesDownload Raw Back to client
1// Copyright (c) Microsoft Corporation.2// Licensed under the MIT License.3 4using System.Collections.Concurrent;5using System.Collections.Generic;6using System.Globalization;7using System.Linq;8using System.Management.Automation.Remoting.Internal;9using System.Threading;10 11using Dbg = System.Management.Automation.Diagnostics;12 13namespace System.Management.Automation14{15    internal abstract class StartableJob : Job16    {17        internal StartableJob(string commandName, string jobName)18            : base(commandName, jobName)19        {20        }21 22        internal abstract void StartJob();23    }24 25    /// <summary>26    /// A job that can throttle execution of child jobs.27    /// </summary>28    internal sealed class ThrottlingJob : Job29    {30        #region IDisposable Members31 32        /// <summary>33        /// Releases resources associated with this object.34        /// </summary>35        protected override void Dispose(bool disposing)36        {37            try38            {39                if (disposing)40                {41                    this.StopJob();42 43                    List<Job> childJobsToDispose;44                    lock (_lockObject)45                    {46                        Dbg.Assert(this.IsFinishedState(this.JobStateInfo.State), "ThrottlingJob should be completed before removing and disposing child jobs");47                        childJobsToDispose = new List<Job>(this.ChildJobs);48                        this.ChildJobs.Clear();49                    }50 51                    foreach (Job childJob in childJobsToDispose)52                    {53                        childJob.Dispose();54                    }55 56                    _jobResultsThrottlingSemaphore?.Dispose();57                    _cancellationTokenSource.Dispose();58                }59            }60            finally61            {62                base.Dispose(disposing);63            }64        }65 66        #endregion67 68        #region Support for progress reporting69 70        private readonly DateTime _progressStartTime = DateTime.UtcNow;71        private readonly int _progressActivityId;72        private readonly object _progressLock = new object();73        private DateTime _progressReportLastTime = DateTime.MinValue;74 75        internal int GetProgressActivityId()76        {77            lock (_progressLock)78            {79                if (_progressReportLastTime.Equals(DateTime.MinValue))80                {81                    try82                    {83                        this.ReportProgress(minimizeFrequentUpdates: false);84                        Dbg.Assert(_progressReportLastTime > DateTime.MinValue, "Progress was reported (lastTimeProgressWasReported)");85                    }86                    catch (PSInvalidOperationException)87                    {88                        // return "no parent activity id" if this ThrottlingJob has already finished89                        return -1;90                    }91                }92 93                return _progressActivityId;94            }95        }96 97        private void ReportProgress(bool minimizeFrequentUpdates)98        {99            lock (_progressLock)100            {101                DateTime now = DateTime.UtcNow;102 103                if (minimizeFrequentUpdates)104                {105                    if ((now - _progressStartTime) < TimeSpan.FromSeconds(1))106                    {107                        return;108                    }109 110                    if ((!_progressReportLastTime.Equals(DateTime.MinValue)) &&111                        (now - _progressReportLastTime < TimeSpan.FromMilliseconds(200)))112                    {113                        return;114                    }115                }116 117                _progressReportLastTime = now;118 119                double workCompleted;120                double totalWork;121                int percentComplete;122                lock (_lockObject)123                {124                    totalWork = _countOfAllChildJobs;125                    workCompleted = this.CountOfFinishedChildJobs;126                }127 128                if (totalWork >= 1.0)129                {130                    percentComplete = (int)(100.0 * workCompleted / totalWork);131                }132                else133                {134                    percentComplete = -1;135                }136 137                percentComplete = Math.Max(-1, Math.Min(100, percentComplete));138 139                var progressRecord = new ProgressRecord(140                    activityId: _progressActivityId,141                    activity: this.Command,142                    statusDescription: this.StatusMessage);143 144                if (this.IsThrottlingJobCompleted)145                {146                    if (_progressReportLastTime.Equals(DateTime.MinValue))147                    {148                        return;149                    }150 151                    progressRecord.RecordType = ProgressRecordType.Completed;152                    progressRecord.PercentComplete = 100;153                    progressRecord.SecondsRemaining = 0;154                }155                else156                {157                    progressRecord.RecordType = ProgressRecordType.Processing;158                    progressRecord.PercentComplete = percentComplete;159                    int? secondsRemaining = null;160                    if (percentComplete >= 0)161                    {162                        secondsRemaining = ProgressRecord.GetSecondsRemaining(_progressStartTime, (double)percentComplete / 100.0);163                    }164 165                    if (secondsRemaining.HasValue)166                    {167                        progressRecord.SecondsRemaining = secondsRemaining.Value;168                    }169                }170 171                this.WriteProgress(progressRecord);172            }173        }174 175        #endregion176 177        /// <summary>178        /// Flags of child jobs of a <see cref="ThrottlingJob"/>179        /// </summary>180        [Flags]181        internal enum ChildJobFlags182        {183            /// <summary>184            /// Child job doesn't have any special properties.185            /// </summary>186            None = 0,187 188            /// <summary>189            /// Child job can call <see cref="ThrottlingJob.AddChildJobWithoutBlocking"/> method190            /// or <see cref="ThrottlingJob.AddChildJobAndPotentiallyBlock(StartableJob, ChildJobFlags)"/>191            /// or <see cref="ThrottlingJob.AddChildJobAndPotentiallyBlock(Cmdlet, StartableJob, ChildJobFlags)"/>192            /// method193            /// of the <see cref="ThrottlingJob"/> instance it belongs to.194            /// </summary>195            CreatesChildJobs = 0x1,196        }197 198        private bool _ownerWontSubmitNewChildJobs = false;199        private readonly HashSet<Guid> _setOfChildJobsThatCanAddMoreChildJobs = new HashSet<Guid>();200 201        private bool IsEndOfChildJobs202        {203            get204            {205                lock (_lockObject)206                {207                    return _isStopping || (_ownerWontSubmitNewChildJobs && _setOfChildJobsThatCanAddMoreChildJobs.Count == 0);208                }209            }210        }211 212        private bool IsThrottlingJobCompleted213        {214            get215            {216                lock (_lockObject)217                {218                    return this.IsEndOfChildJobs && (_countOfAllChildJobs <= this.CountOfFinishedChildJobs);219                }220            }221        }222 223        private readonly bool _cmdletMode;224 225        private int _countOfAllChildJobs;226 227        private int _countOfBlockedChildJobs;228 229        private int _countOfFailedChildJobs;230        private int _countOfStoppedChildJobs;231        private int _countOfSuccessfullyCompletedChildJobs;232 233        private int CountOfFinishedChildJobs234        {235            get236            {237                lock (_lockObject)238                {239                    return _countOfFailedChildJobs + _countOfStoppedChildJobs + _countOfSuccessfullyCompletedChildJobs;240                }241            }242        }243 244        private int CountOfRunningOrReadyToRunChildJobs245        {246            get247            {248                lock (_lockObject)249                {250                    return _countOfAllChildJobs - this.CountOfFinishedChildJobs;251                }252            }253        }254 255        private readonly object _lockObject = new object();256 257        /// <summary>258        /// Creates a new <see cref="ThrottlingJob"/> object.259        /// </summary>260        /// <param name="command">Command invoked by this job object.</param>261        /// <param name="jobName">Friendly name for the job object.</param>262        /// <param name="jobTypeName">Name describing job type.</param>263        /// <param name="maximumConcurrentChildJobs">264        /// The maximum number of child jobs that can be running at any given point in time.265        /// Passing 0 requests to turn off throttling (i.e. allow unlimited number of child jobs to run)266        /// </param>267        /// <param name="cmdletMode">268        /// <see langword="true"/> if this <see cref="ThrottlingJob"/> is used from a cmdlet invoked without -AsJob switch.269        /// <see langword="false"/> if this <see cref="ThrottlingJob"/> is used from a cmdlet invoked with -AsJob switch.270        ///271        /// If <paramref name="cmdletMode"/> is <see langword="true"/>, then272        /// memory can be managed more aggressively (for example ChildJobs can be discarded as soon as they complete)273        /// because the <see cref="ThrottlingJob"/> is not exposed to the end user.274        /// </param>275        internal ThrottlingJob(string command, string jobName, string jobTypeName, int maximumConcurrentChildJobs, bool cmdletMode)276            : base(command, jobName)277        {278            this.Results.BlockingEnumerator = true;279            _cmdletMode = cmdletMode;280            this.PSJobTypeName = jobTypeName;281            if (_cmdletMode)282            {283                _jobResultsThrottlingSemaphore = new SemaphoreSlim(ForwardingHelper.AggregationQueueMaxCapacity);284            }285 286            _progressActivityId = new Random(this.GetHashCode()).Next();287 288            this.SetupThrottlingQueue(maximumConcurrentChildJobs);289        }290 291        internal void AddChildJobAndPotentiallyBlock(292            StartableJob childJob,293            ChildJobFlags flags)294        {295            using (var jobGotEnqueued = new ManualResetEventSlim(initialState: false))296            {297                ArgumentNullException.ThrowIfNull(childJob);298 299                this.AddChildJobWithoutBlocking(childJob, flags, jobGotEnqueued.Set);300                jobGotEnqueued.Wait();301            }302        }303 304        internal void AddChildJobAndPotentiallyBlock(305            Cmdlet cmdlet,306            StartableJob childJob,307            ChildJobFlags flags)308        {309            using (var forwardingCancellation = new CancellationTokenSource())310            {311                ArgumentNullException.ThrowIfNull(childJob);312 313                this.AddChildJobWithoutBlocking(childJob, flags, forwardingCancellation.Cancel);314                this.ForwardAllResultsToCmdlet(cmdlet, forwardingCancellation.Token);315            }316        }317 318        private bool _alreadyDisabledFlowControlForPendingJobsQueue = false;319 320        internal void DisableFlowControlForPendingJobsQueue()321        {322            if (!_cmdletMode || _alreadyDisabledFlowControlForPendingJobsQueue)323            {324                return;325            }326 327            _alreadyDisabledFlowControlForPendingJobsQueue = true;328 329            lock (_lockObject)330            {331                _maxReadyToRunJobs = int.MaxValue;332 333                while (_actionsForUnblockingChildAdditions.Count > 0)334                {335                    Action a = _actionsForUnblockingChildAdditions.Dequeue();336                    a?.Invoke();337                }338            }339        }340 341        private bool _alreadyDisabledFlowControlForPendingCmdletActionsQueue = false;342 343        internal void DisableFlowControlForPendingCmdletActionsQueue()344        {345            if (!_cmdletMode || _alreadyDisabledFlowControlForPendingCmdletActionsQueue)346            {347                return;348            }349 350            _alreadyDisabledFlowControlForPendingCmdletActionsQueue = true;351 352            long slotsToRelease = (long)(int.MaxValue / 2) - (long)(_jobResultsThrottlingSemaphore.CurrentCount);353            if ((slotsToRelease > 0) && (slotsToRelease < int.MaxValue))354            {355                _jobResultsThrottlingSemaphore.Release((int)slotsToRelease);356            }357        }358 359        /// <summary>360        /// Adds and starts a child job.361        /// </summary>362        /// <param name="childJob">Child job to add.</param>363        /// <param name="flags">Flags of the child job.</param>364        /// <param name="jobEnqueuedAction">Action to run after enqueuing the job.</param>365        /// <exception cref="ArgumentException">366        /// Thrown when the child job is not in the <see cref="JobState.NotStarted"/> state.367        /// (because this can lead to race conditions - the child job can finish before the parent job has a chance to register for child job events)368        /// </exception>369        internal void AddChildJobWithoutBlocking(StartableJob childJob, ChildJobFlags flags, Action jobEnqueuedAction = null)370        {371            ArgumentNullException.ThrowIfNull(childJob);372            if (childJob.JobStateInfo.State != JobState.NotStarted)373            {374                throw new ArgumentException(RemotingErrorIdStrings.ThrottlingJobChildAlreadyRunning, nameof(childJob));375            }376                377            this.AssertNotDisposed();378 379            JobStateInfo newJobStateInfo = null;380            lock (_lockObject)381            {382                if (this.IsEndOfChildJobs)383                {384                    throw new InvalidOperationException(RemotingErrorIdStrings.ThrottlingJobChildAddedAfterEndOfChildJobs);385                }386 387                if (_isStopping)388                {389                    return;390                }391 392                if (_countOfAllChildJobs == 0)393                {394                    newJobStateInfo = new JobStateInfo(JobState.Running);395                }396 397                if ((ChildJobFlags.CreatesChildJobs & flags) == ChildJobFlags.CreatesChildJobs)398                {399                    _setOfChildJobsThatCanAddMoreChildJobs.Add(childJob.InstanceId);400                }401 402                this.ChildJobs.Add(childJob);403                _childJobLocations.Add(childJob.Location);404                _countOfAllChildJobs++;405 406                this.WriteWarningAboutHighUsageOfFlowControlBuffers(this.CountOfRunningOrReadyToRunChildJobs);407                if (this.CountOfRunningOrReadyToRunChildJobs > _maxReadyToRunJobs)408                {409                    _actionsForUnblockingChildAdditions.Enqueue(jobEnqueuedAction);410                }411                else412                {413                    jobEnqueuedAction?.Invoke();414                }415            }416 417            if (newJobStateInfo != null)418            {419                this.SetJobState(newJobStateInfo.State, newJobStateInfo.Reason);420            }421 422            this.ChildJobAdded.SafeInvoke(this, new ThrottlingJobChildAddedEventArgs(childJob));423 424            childJob.SetParentActivityIdGetter(this.GetProgressActivityId);425            childJob.StateChanged += this.childJob_StateChanged;426            if (_cmdletMode)427            {428                childJob.Results.DataAdded += childJob_ResultsAdded;429            }430 431            this.EnqueueReadyToRunChildJob(childJob);432 433            this.ReportProgress(minimizeFrequentUpdates: true);434        }435 436        private void childJob_ResultsAdded(object sender, DataAddedEventArgs e)437        {438            Dbg.Assert(_jobResultsThrottlingSemaphore != null, "JobResultsThrottlingSemaphore should be non-null if childJob_ResultsAdded handled is registered");439            try440            {441                long jobResultsUpdatedCount = Interlocked.Increment(ref _jobResultsCurrentCount);442                this.WriteWarningAboutHighUsageOfFlowControlBuffers(jobResultsUpdatedCount);443 444                _jobResultsThrottlingSemaphore.Wait(_cancellationTokenSource.Token);445            }446            catch (ObjectDisposedException)447            {448            }449            catch (OperationCanceledException)450            {451            }452        }453 454        private readonly object _alreadyWroteFlowControlBuffersHighMemoryUsageWarningLock = new object();455        private bool _alreadyWroteFlowControlBuffersHighMemoryUsageWarning;456 457        private const long FlowControlBuffersHighMemoryUsageThreshold = 30000;458 459        private void WriteWarningAboutHighUsageOfFlowControlBuffers(long currentCount)460        {461            if (!_cmdletMode)462            {463                return;464            }465 466            if (currentCount < FlowControlBuffersHighMemoryUsageThreshold)467            {468                return;469            }470 471            lock (_alreadyWroteFlowControlBuffersHighMemoryUsageWarningLock)472            {473                if (_alreadyWroteFlowControlBuffersHighMemoryUsageWarning)474                {475                    return;476                }477 478                _alreadyWroteFlowControlBuffersHighMemoryUsageWarning = true;479            }480 481            string warningMessage = string.Format(482                CultureInfo.InvariantCulture,483                RemotingErrorIdStrings.ThrottlingJobFlowControlMemoryWarning,484                this.Command);485            this.WriteWarning(warningMessage);486        }487 488        internal event EventHandler<ThrottlingJobChildAddedEventArgs> ChildJobAdded;489 490        private int _maximumConcurrentChildJobs;491        private int _extraCapacityForRunningQueryJobs;492        private int _extraCapacityForRunningAllJobs;493        private bool _inBoostModeToPreventQueryJobDeadlock;494        private Queue<StartableJob> _readyToRunQueryJobs;495        private Queue<StartableJob> _readyToRunRegularJobs;496        private Queue<Action> _actionsForUnblockingChildAdditions;497        private int _maxReadyToRunJobs;498        private readonly SemaphoreSlim _jobResultsThrottlingSemaphore;499        private long _jobResultsCurrentCount;500 501        private static readonly int s_maximumReadyToRunJobs = 10000;502 503        private void SetupThrottlingQueue(int maximumConcurrentChildJobs)504        {505            _maximumConcurrentChildJobs = maximumConcurrentChildJobs > 0 ? maximumConcurrentChildJobs : int.MaxValue;506            if (_cmdletMode)507            {508                _maxReadyToRunJobs = s_maximumReadyToRunJobs;509            }510            else511            {512                _maxReadyToRunJobs = int.MaxValue;513            }514 515            _extraCapacityForRunningAllJobs = _maximumConcurrentChildJobs;516            _extraCapacityForRunningQueryJobs = Math.Max(1, _extraCapacityForRunningAllJobs / 2);517 518            _inBoostModeToPreventQueryJobDeadlock = false;519            _readyToRunQueryJobs = new Queue<StartableJob>();520            _readyToRunRegularJobs = new Queue<StartableJob>();521            _actionsForUnblockingChildAdditions = new Queue<Action>();522        }523 524        private void StartChildJobIfPossible()525        {526            StartableJob readyToRunChildJob = null;527            lock (_lockObject)528            {529                do530                {531                    if ((_readyToRunQueryJobs.Count > 0) &&532                        (_extraCapacityForRunningQueryJobs > 0) &&533                        (_extraCapacityForRunningAllJobs > 0))534                    {535                        _extraCapacityForRunningQueryJobs--;536                        _extraCapacityForRunningAllJobs--;537                        readyToRunChildJob = _readyToRunQueryJobs.Dequeue();538                        break;539                    }540 541                    if ((_readyToRunRegularJobs.Count > 0) &&542                        (_extraCapacityForRunningAllJobs > 0))543                    {544                        _extraCapacityForRunningAllJobs--;545                        readyToRunChildJob = _readyToRunRegularJobs.Dequeue();546                        break;547                    }548                } while (false);549            }550 551            readyToRunChildJob?.StartJob();552        }553 554        private void EnqueueReadyToRunChildJob(StartableJob childJob)555        {556            lock (_lockObject)557            {558                bool isQueryJob = _setOfChildJobsThatCanAddMoreChildJobs.Contains(childJob.InstanceId);559                if (isQueryJob &&560                    !_inBoostModeToPreventQueryJobDeadlock &&561                    (_maximumConcurrentChildJobs == 1))562                {563                    _inBoostModeToPreventQueryJobDeadlock = true;564                    _extraCapacityForRunningAllJobs++;565                }566 567                if (isQueryJob)568                {569                    _readyToRunQueryJobs.Enqueue(childJob);570                }571                else572                {573                    _readyToRunRegularJobs.Enqueue(childJob);574                }575            }576 577            StartChildJobIfPossible();578        }579 580        private void MakeRoomForRunningOtherJobs(Job completedChildJob)581        {582            lock (_lockObject)583            {584                _extraCapacityForRunningAllJobs++;585 586                bool isQueryJob = _setOfChildJobsThatCanAddMoreChildJobs.Contains(completedChildJob.InstanceId);587                if (isQueryJob)588                {589                    _setOfChildJobsThatCanAddMoreChildJobs.Remove(completedChildJob.InstanceId);590 591                    _extraCapacityForRunningQueryJobs++;592                    if (_inBoostModeToPreventQueryJobDeadlock && (_setOfChildJobsThatCanAddMoreChildJobs.Count == 0))593                    {594                        _inBoostModeToPreventQueryJobDeadlock = false;595                        _extraCapacityForRunningAllJobs--;596                    }597                }598            }599 600            StartChildJobIfPossible();601        }602 603        private void FigureOutIfThrottlingJobIsCompleted()604        {605            JobStateInfo finalJobStateInfo = null;606            lock (_lockObject)607            {608                if (this.IsThrottlingJobCompleted && !IsFinishedState(this.JobStateInfo.State))609                {610                    if (_isStopping)611                    {612                        finalJobStateInfo = new JobStateInfo(JobState.Stopped, null);613                    }614                    else if (_countOfFailedChildJobs > 0)615                    {616                        finalJobStateInfo = new JobStateInfo(JobState.Failed, null);617                    }618                    else if (_countOfStoppedChildJobs > 0)619                    {620                        finalJobStateInfo = new JobStateInfo(JobState.Stopped, null);621                    }622                    else623                    {624                        finalJobStateInfo = new JobStateInfo(JobState.Completed);625                    }626                }627            }628 629            if (finalJobStateInfo != null)630            {631                this.SetJobState(finalJobStateInfo.State, finalJobStateInfo.Reason);632                this.CloseAllStreams();633            }634        }635 636        /// <summary>637        /// Notifies this <see cref="ThrottlingJob"/> object that no more child jobs will be added.638        /// </summary>639        internal void EndOfChildJobs()640        {641            this.AssertNotDisposed();642            lock (_lockObject)643            {644                _ownerWontSubmitNewChildJobs = true;645            }646 647            this.FigureOutIfThrottlingJobIsCompleted();648        }649 650        /// <summary>651        /// Stop this job object and all the <see cref="System.Management.Automation.Job.ChildJobs"/>.652        /// </summary>653        public override void StopJob()654        {655            List<Job> childJobsToStop = null;656            lock (_lockObject)657            {658                if (!(_isStopping || this.IsThrottlingJobCompleted))659                {660                    _isStopping = true;661                    childJobsToStop = this.GetChildJobsSnapshot();662                }663            }664 665            if (childJobsToStop != null)666            {667                this.SetJobState(JobState.Stopping);668 669                _cancellationTokenSource.Cancel();670                foreach (Job childJob in childJobsToStop)671                {672                    if (!childJob.IsFinishedState(childJob.JobStateInfo.State))673                    {674                        childJob.StopJob();675                    }676                }677 678                this.FigureOutIfThrottlingJobIsCompleted();679            }680 681            this.Finished.WaitOne();682        }683 684        private bool _isStopping;685        private readonly CancellationTokenSource _cancellationTokenSource = new CancellationTokenSource();686 687        private void childJob_StateChanged(object sender, JobStateEventArgs e)688        {689            Dbg.Assert(sender != null, "Only our internal implementation of Job should raise this event and it should make sure that sender != null");690            Dbg.Assert(sender is Job, "Only our internal implementation of Job should raise this event and it should make sure that sender is Job");691            var childJob = (Job)sender;692 693            if ((e.PreviousJobStateInfo.State == JobState.Blocked) && (e.JobStateInfo.State != JobState.Blocked))694            {695                bool parentJobGotUnblocked = false;696                lock (_lockObject)697                {698                    _countOfBlockedChildJobs--;699                    if (_countOfBlockedChildJobs == 0)700                    {701                        parentJobGotUnblocked = true;702                    }703                }704 705                if (parentJobGotUnblocked)706                {707                    this.SetJobState(JobState.Running);708                }709            }710 711            switch (e.JobStateInfo.State)712            {713                // intermediate states714                case JobState.Blocked:715                    lock (_lockObject)716                    {717                        _countOfBlockedChildJobs++;718                    }719 720                    this.SetJobState(JobState.Blocked);721                    break;722 723                // 3 finished states724                case JobState.Failed:725                case JobState.Stopped:726                case JobState.Completed:727                    childJob.StateChanged -= childJob_StateChanged;728                    this.MakeRoomForRunningOtherJobs(childJob);729                    lock (_lockObject)730                    {731                        if (e.JobStateInfo.State == JobState.Failed)732                        {733                            _countOfFailedChildJobs++;734                        }735                        else if (e.JobStateInfo.State == JobState.Stopped)736                        {737                            _countOfStoppedChildJobs++;738                        }739                        else if (e.JobStateInfo.State == JobState.Completed)740                        {741                            _countOfSuccessfullyCompletedChildJobs++;742                        }743 744                        if (_actionsForUnblockingChildAdditions.Count > 0)745                        {746                            Action a = _actionsForUnblockingChildAdditions.Dequeue();747                            a?.Invoke();748                        }749 750                        if (_cmdletMode)751                        {752                            foreach (PSStreamObject streamObject in childJob.Results.ReadAll())753                            {754                                this.Results.Add(streamObject);755                            }756 757                            this.ChildJobs.Remove(childJob);758                            _setOfChildJobsThatCanAddMoreChildJobs.Remove(childJob.InstanceId);759                            childJob.Dispose();760                        }761                    }762 763                    this.ReportProgress(minimizeFrequentUpdates: !this.IsThrottlingJobCompleted);764                    break;765 766                default:767                    // do nothing768                    break;769            }770 771            this.FigureOutIfThrottlingJobIsCompleted();772        }773 774        private List<Job> GetChildJobsSnapshot()775        {776            lock (_lockObject)777            {778                return new List<Job>(this.ChildJobs);779            }780        }781 782        /// <summary>783        /// Indicates if job has more data available.784        /// <see langword="true"/> if any of the child jobs have more data OR if <see cref="EndOfChildJobs"/> have not been called yet;785        /// <see langword="false"/> otherwise.786        /// </summary>787        public override bool HasMoreData788        {789            get790            {791                return this.GetChildJobsSnapshot().Any(static childJob => childJob.HasMoreData) || (this.Results.Count != 0);792            }793        }794 795        /// <summary>796        /// Comma-separated list of locations of <see cref="System.Management.Automation.Job.ChildJobs"/>.797        /// </summary>798        public override string Location799        {800            get801            {802                lock (_lockObject)803                {804                    return string.Join(", ", _childJobLocations);805                }806            }807        }808 809        private readonly HashSet<string> _childJobLocations = new HashSet<string>(StringComparer.OrdinalIgnoreCase);810 811        /// <summary>812        /// Status message associated with the Job.813        /// </summary>814        public override string StatusMessage815        {816            get817            {818                int completedChildJobs;819                int totalChildJobs;820                lock (_lockObject)821                {822                    completedChildJobs = this.CountOfFinishedChildJobs;823                    totalChildJobs = _countOfAllChildJobs;824                }825 826                string totalChildJobsString = totalChildJobs.ToString(CultureInfo.CurrentCulture);827                if (!this.IsEndOfChildJobs)828                {829                    totalChildJobsString += "+";830                }831 832                return string.Format(833                    CultureInfo.CurrentUICulture,834                    RemotingErrorIdStrings.ThrottlingJobStatusMessage,835                    completedChildJobs,836                    totalChildJobsString);837            }838        }839 840        #region Forwarding results to a cmdlet841 842        internal override void ForwardAvailableResultsToCmdlet(Cmdlet cmdlet)843        {844            this.AssertNotDisposed();845 846            base.ForwardAvailableResultsToCmdlet(cmdlet);847            foreach (Job childJob in this.GetChildJobsSnapshot())848            {849                childJob.ForwardAvailableResultsToCmdlet(cmdlet);850            }851        }852 853        private sealed class ForwardingHelper : IDisposable854        {855            // This is higher than 1000 used in856            //      RxExtensionMethods+ToEnumerableObserver<T>.BlockingCollectionCapacity857            // and in858            //      RemoteDiscoveryHelper.BlockingCollectionCapacity859            // It needs to be higher, because the high value is used as an attempt to workaround the fact that860            // WSMan will timeout if an OnNext call blocks for more than X minutes.861 862            // This is a static field (instead of a constant) to make it possible to set through tests (and/or by customers if needed for a workaround)863            internal static readonly int AggregationQueueMaxCapacity = 10000;864 865            private readonly ThrottlingJob _throttlingJob;866 867            private readonly object _myLock;868            private readonly BlockingCollection<PSStreamObject> _aggregatedResults;869            private readonly HashSet<Job> _monitoredJobs;870 871            private readonly CancellationTokenSource _cancellationTokenSource = new CancellationTokenSource();872            private bool _disposed;873 874            private ForwardingHelper(ThrottlingJob throttlingJob)875            {876                _throttlingJob = throttlingJob;877 878                _myLock = new object();879                _monitoredJobs = new HashSet<Job>();880 881                _aggregatedResults = new BlockingCollection<PSStreamObject>();882            }883 884            private void StartMonitoringJob(Job job)885            {886                lock (_myLock)887                {888                    if (_disposed || _stoppedMonitoringAllJobs)889                    {890                        return;891                    }892 893                    if (_monitoredJobs.Contains(job))894                    {895                        return;896                    }897 898                    _monitoredJobs.Add(job);899 900                    job.Results.DataAdded += this.MonitoredJobResults_DataAdded;901                    job.StateChanged += MonitoredJob_StateChanged;902                }903 904                this.AggregateJobResults(job.Results);905                this.CheckIfMonitoredJobIsComplete(job);906            }907 908            private void StopMonitoringJob(Job job)909            {910                lock (_myLock)911                {912                    if (_monitoredJobs.Contains(job))913                    {914                        job.Results.DataAdded -= this.MonitoredJobResults_DataAdded;915                        job.StateChanged -= this.MonitoredJob_StateChanged;916                        _monitoredJobs.Remove(job);917                    }918                }919            }920 921            private void AggregateJobResults(PSDataCollection<PSStreamObject> resultsCollection)922            {923                lock (_myLock)924                {925                    // try not to remove results from a job, unless it seems safe ...926                    if (_disposed || _stoppedMonitoringAllJobs || _aggregatedResults.IsAddingCompleted || _cancellationTokenSource.IsCancellationRequested)927                    {928                        return;929                    }930                }931 932                // ... and after removing the results via ReadAll, we have to make sure that we don't drop them ...933                foreach (var result in resultsCollection.ReadAll())934                {935                    bool successfullyAggregatedResult = false;936                    try937                    {938                        lock (_myLock)939                        {940                            // try not to remove results from a job, unless it seems safe ...941                            if (!(_disposed || _stoppedMonitoringAllJobs || _aggregatedResults.IsAddingCompleted || _cancellationTokenSource.IsCancellationRequested))942                            {943                                _aggregatedResults.Add(result, _cancellationTokenSource.Token);944                                successfullyAggregatedResult = true;945                            }946                        }947                    }948                    catch (Exception) // BlockingCollection.Add can throw undocumented exceptions - we cannot just catch InvalidOperationException949                    {950                    }951 952                    // ... so if _aggregatedResults is not accepting new results, we will store them in the throttling job953                    if (!successfullyAggregatedResult)954                    {955                        this.StopMonitoringJob(_throttlingJob);956                        try957                        {958                            _throttlingJob.Results.Add(result);959                        }960                        catch (InvalidOperationException)961                        {962                            Dbg.Assert(false, "ThrottlingJob.Results was already closed when trying to preserve results aggregated by ForwardingHelper");963                        }964                    }965                }966            }967 968            private void CancelForwarding()969            {970                _cancellationTokenSource.Cancel();971                lock (_myLock)972                {973                    Dbg.Assert(!_disposed, "CancelForwarding should be unregistered before ForwardingHelper gets disposed");974                    _aggregatedResults.CompleteAdding();975                }976            }977 978            private void CheckIfMonitoredJobIsComplete(Job job)979            {980                CheckIfMonitoredJobIsComplete(job, job.JobStateInfo.State);981            }982 983            private void CheckIfMonitoredJobIsComplete(Job job, JobState jobState)984            {985                if (job.IsFinishedState(jobState))986                {987                    lock (_myLock)988                    {989                        this.StopMonitoringJob(job);990                    }991                }992            }993 994            private void CheckIfThrottlingJobIsComplete()995            {996                if (_throttlingJob.IsThrottlingJobCompleted)997                {998                    List<PSDataCollection<PSStreamObject>> resultsToAggregate = new List<PSDataCollection<PSStreamObject>>();999                    lock (_myLock)1000                    {1001                        foreach (Job registeredJob in _monitoredJobs)1002                        {1003                            resultsToAggregate.Add(registeredJob.Results);1004                        }1005 1006                        foreach (Job throttledJob in _throttlingJob.GetChildJobsSnapshot())1007                        {1008                            resultsToAggregate.Add(throttledJob.Results);1009                        }1010 1011                        resultsToAggregate.Add(_throttlingJob.Results);1012                    }1013 1014                    foreach (PSDataCollection<PSStreamObject> resultToAggregate in resultsToAggregate)1015                    {1016                        this.AggregateJobResults(resultToAggregate);1017                    }1018 1019                    lock (_myLock)1020                    {1021                        if (!_disposed && !_aggregatedResults.IsAddingCompleted)1022                        {1023                            _aggregatedResults.CompleteAdding();1024                        }1025                    }1026                }1027            }1028 1029            private void MonitoredJobResults_DataAdded(object sender, DataAddedEventArgs e)1030            {1031                var resultsCollection = (PSDataCollection<PSStreamObject>)sender;1032                this.AggregateJobResults(resultsCollection);1033            }1034 1035            private void MonitoredJob_StateChanged(object sender, JobStateEventArgs e)1036            {1037                var job = (Job)sender;1038                this.CheckIfMonitoredJobIsComplete(job, e.JobStateInfo.State);1039            }1040 1041            private void ThrottlingJob_ChildJobAdded(object sender, ThrottlingJobChildAddedEventArgs e)1042            {1043                this.StartMonitoringJob(e.AddedChildJob);1044            }1045 1046            private void ThrottlingJob_StateChanged(object sender, JobStateEventArgs e)1047            {1048                this.CheckIfThrottlingJobIsComplete();1049            }1050 1051            private void AttemptToPreserveAggregatedResults()1052            {1053#if DEBUG1054                lock (_myLock)1055                {1056                    Dbg.Assert(!_disposed, "AttemptToPreserveAggregatedResults should be called before disposing ForwardingHelper");1057                    Dbg.Assert(_stoppedMonitoringAllJobs, "Caller should guarantee no-more-results before calling AttemptToPreserveAggregatedResults (1)");1058                    Dbg.Assert(_aggregatedResults.IsAddingCompleted, "Caller should guarantee no-more-results before calling AttemptToPreserveAggregatedResults (2)");1059                }1060#endif1061 1062                bool isThrottlingJobFinished = false;1063                foreach (var aggregatedButNotYetProcessedResult in _aggregatedResults)1064                {1065                    if (!isThrottlingJobFinished)1066                    {1067                        try1068                        {1069                            _throttlingJob.Results.Add(aggregatedButNotYetProcessedResult);1070                        }1071                        catch (PSInvalidOperationException)1072                        {1073                            isThrottlingJobFinished = _throttlingJob.IsFinishedState(_throttlingJob.JobStateInfo.State);1074                            Dbg.Assert(isThrottlingJobFinished, "Buffers should not be closed before throttling job is stopped");1075                        }1076                    }1077                }1078            }1079 1080#if DEBUG1081            // CDXML_CLIXML_TEST testability hook1082 1083            private static readonly bool s_isCliXmlTestabilityHookActive = GetIsCliXmlTestabilityHookActive();1084 1085            private static bool GetIsCliXmlTestabilityHookActive()1086            {1087                return !string.IsNullOrEmpty(Environment.GetEnvironmentVariable("CDXML_CLIXML_TEST"));1088            }1089 1090            internal static void ProcessCliXmlTestabilityHook(PSStreamObject streamObject)1091            {1092                if (!s_isCliXmlTestabilityHookActive)1093                {1094                    return;1095                }1096 1097                if (streamObject.ObjectType != PSStreamObjectType.Output)1098                {1099                    return;1100                }1101 1102                if (streamObject.Value == null)1103                {1104                    return;1105                }1106 1107                if (!(PSObject.AsPSObject(streamObject.Value).BaseObject.GetType().Name.Equals("CimInstance")))1108                {1109                    return;1110                }1111 1112                string serializedForm = PSSerializer.Serialize(streamObject.Value, depth: 1);1113                object deserializedObject = PSSerializer.Deserialize(serializedForm);1114                streamObject.Value = PSObject.AsPSObject(deserializedObject).BaseObject;1115            }1116#endif1117 1118            private void ForwardResults(Cmdlet cmdlet)1119            {1120                try1121                {1122                    foreach (var result in _aggregatedResults.GetConsumingEnumerable(_throttlingJob._cancellationTokenSource.Token))1123                    {1124                        if (result != null)1125                        {1126#if DEBUG1127                            // CDXML_CLIXML_TEST testability hook1128                            ProcessCliXmlTestabilityHook(result);1129#endif1130                            try1131                            {1132                                result.WriteStreamObject(cmdlet);1133                            }1134                            finally1135                            {1136                                if (_throttlingJob._cmdletMode)1137                                {1138                                    Dbg.Assert(_throttlingJob._jobResultsThrottlingSemaphore != null, "JobResultsThrottlingSemaphore should be present in cmdlet mode");1139                                    Interlocked.Decrement(ref _throttlingJob._jobResultsCurrentCount);1140                                    _throttlingJob._jobResultsThrottlingSemaphore.Release();1141                                }1142                            }1143                        }1144                    }1145                }1146                catch1147                {1148                    this.StopMonitoringAllJobs();1149                    this.AttemptToPreserveAggregatedResults();1150 1151                    throw;1152                }1153            }1154 1155            private bool _stoppedMonitoringAllJobs;1156 1157            private void StopMonitoringAllJobs()1158            {1159                _cancellationTokenSource.Cancel();1160                lock (_myLock)1161                {1162                    _stoppedMonitoringAllJobs = true;1163 1164                    List<Job> snapshotOfCurrentlyMonitoredJobs = _monitoredJobs.ToList();1165                    foreach (Job monitoredJob in snapshotOfCurrentlyMonitoredJobs)1166                    {1167                        this.StopMonitoringJob(monitoredJob);1168                    }1169 1170                    Dbg.Assert(_monitoredJobs.Count == 0, "No monitored jobs should be left after ForwardingHelper is disposed");1171 1172                    if (!_disposed && !_aggregatedResults.IsAddingCompleted)1173                    {1174                        _aggregatedResults.CompleteAdding();1175                    }1176                }1177            }1178 1179            public void Dispose()1180            {1181                GC.SuppressFinalize(this);1182                _cancellationTokenSource.Cancel();1183                lock (_myLock)1184                {1185                    if (_disposed)1186                    {1187                        return;1188                    }1189 1190                    this.StopMonitoringAllJobs();1191                    _aggregatedResults.Dispose();1192                    _cancellationTokenSource.Dispose();1193                    _disposed = true;1194                }1195            }1196 1197            public static void ForwardAllResultsToCmdlet(ThrottlingJob throttlingJob, Cmdlet cmdlet, CancellationToken? cancellationToken)1198            {1199                using (var helper = new ForwardingHelper(throttlingJob))1200                {

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