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