MegaBites-AI/Windows-powershell
0372
1// Copyright (c) Microsoft Corporation.2// Licensed under the MIT License.3 4using System;5using System.Collections;6using System.Collections.Generic;7using System.Collections.ObjectModel;8using System.Diagnostics.CodeAnalysis;9using System.Management.Automation;10using System.Management.Automation.Internal;11using System.Management.Automation.Remoting;12using System.Management.Automation.Remoting.Internal;13using System.Management.Automation.Runspaces;14using System.Management.Automation.Tracing;15using System.Threading;16 17using Dbg = System.Management.Automation.Diagnostics;18 19// Stops compiler from warning about unknown warnings20#pragma warning disable 1634, 169121 22namespace Microsoft.PowerShell.Commands23{24 /// <summary>25 /// Cmdlet used for receiving results from job object.26 /// This cmdlet is intended to have a slightly different behavior27 /// in the following two cases:28 /// 1. The job object to receive results from is a PSRemotingJob29 /// In this case, the cmdlet can use two additional30 /// parameters to filter results - ComputerName and Runspace31 /// The parameters help filter out results for a specified32 /// computer or runspace from the job object33 ///34 /// $job = Start-PSJob -Command 'get-process' -ComputerName server1, server235 /// Receive-PSJob -Job $job -ComputerName server136 ///37 /// $job = Start-PSJob -Command 'get-process' -Session $r1, $r238 /// Receive-PSJob -Job $job -Session $r139 ///40 /// 2. The job object to receive results is a PSJob (or derivative41 /// other than PSRemotingJob)42 /// In this case, the user cannot will use the location parameter43 /// to do any filtering and will not have ComputerName and Runspace44 /// parameters45 ///46 /// $job = Get-WMIObject '....' -AsJob47 /// Receive-PSJob -Job $job -Location "Server2"48 ///49 /// The following will result in an error:50 ///51 /// $job = Get-WMIObject '....' -AsJob52 /// Receive-PSJob -Job $job -ComputerName "Server2"53 /// The parameter ComputerName cannot be used with jobs which are54 /// not PSRemotingJob.55 /// </summary>56 [Cmdlet(VerbsCommunications.Receive, "Job", DefaultParameterSetName = ReceiveJobCommand.LocationParameterSet,57 HelpUri = "https://go.microsoft.com/fwlink/?LinkID=2096965", RemotingCapability = RemotingCapability.SupportedByCommand)]58 public class ReceiveJobCommand : JobCmdletBase, IDisposable59 {60 #region Properties61 62 /// <summary>63 /// Job object from which specific results need to64 /// be extracted.65 /// </summary>66 [Parameter(Position = 0,67 Mandatory = true,68 ValueFromPipeline = true,69 ValueFromPipelineByPropertyName = true,70 ParameterSetName = ReceiveJobCommand.ComputerNameParameterSet)]71 [Parameter(Position = 0,72 Mandatory = true,73 ValueFromPipeline = true,74 ValueFromPipelineByPropertyName = true,75 ParameterSetName = ReceiveJobCommand.SessionParameterSet)]76 [Parameter(Position = 0,77 Mandatory = true,78 ValueFromPipeline = true,79 ValueFromPipelineByPropertyName = true,80 ParameterSetName = ReceiveJobCommand.LocationParameterSet)]81 [SuppressMessage("Microsoft.Performance", "CA1819:PropertiesShouldNotReturnArrays")]82 public Job[] Job83 {84 get85 {86 return _jobs;87 }88 89 set90 {91 _jobs = value;92 }93 }94 95 private Job[] _jobs;96 97 /// <summary>98 /// Name of the computer for which the results needs to be99 /// returned.100 /// </summary>101 [Parameter(ValueFromPipelineByPropertyName = true,102 ParameterSetName = ReceiveJobCommand.ComputerNameParameterSet,103 Position = 1)]104 [Alias("Cn")]105 [SuppressMessage("Microsoft.Performance", "CA1819:PropertiesShouldNotReturnArrays")]106 [ValidateNotNullOrEmpty]107 public string[] ComputerName108 {109 get110 {111 return _computerNames;112 }113 114 set115 {116 _computerNames = value;117 }118 }119 120 private string[] _computerNames;121 122 /// <summary>123 /// Locations for which the results needs to be returned.124 /// This will cater to all kinds of jobs and not only125 /// remoting jobs.126 /// </summary>127 [Parameter(ParameterSetName = ReceiveJobCommand.LocationParameterSet,128 Position = 1)]129 [ValidateNotNullOrEmpty]130 [SuppressMessage("Microsoft.Performance", "CA1819:PropertiesShouldNotReturnArrays")]131 public string[] Location132 {133 get134 {135 return _locations;136 }137 138 set139 {140 _locations = value;141 }142 }143 144 private string[] _locations;145 146 /// <summary>147 /// Runspaces for which the results needs to be148 /// returned.149 /// </summary>150 [Parameter(ValueFromPipelineByPropertyName = true,151 ParameterSetName = ReceiveJobCommand.SessionParameterSet,152 Position = 1)]153 [ValidateNotNull]154 [SuppressMessage("Microsoft.Performance", "CA1819:PropertiesShouldNotReturnArrays")]155 public PSSession[] Session156 {157 get158 {159 return _remoteRunspaceInfos;160 }161 162 set163 {164 _remoteRunspaceInfos = value;165 }166 }167 168 private PSSession[] _remoteRunspaceInfos;169 170 /// <summary>171 /// If the results need to be not removed from the store172 /// after being written. Default is results are removed.173 /// </summary>174 [Parameter]175 public SwitchParameter Keep176 {177 get178 {179 return !_flush;180 }181 182 set183 {184 _flush = !value;185 ValidateWait();186 }187 }188 189 private bool _flush = true;190 191 /// <summary>192 /// </summary>193 [Parameter]194 public SwitchParameter NoRecurse195 {196 get197 {198 return !_recurse;199 }200 201 set202 {203 _recurse = !value;204 }205 }206 207 private bool _recurse = true;208 209 /// <summary>210 /// </summary>211 [Parameter]212 public SwitchParameter Force213 { get; set; }214 215 /// <summary>216 /// </summary>217 public override JobState State218 {219 get220 {221 return JobState.NotStarted;222 }223 }224 225 /// <summary>226 /// </summary>227 public override Hashtable Filter228 {229 get { return null; }230 }231 232 /// <summary>233 /// </summary>234 public override string[] Command235 {236 get237 {238 return null;239 }240 }241 242 /// <summary>243 /// </summary>244 protected const string LocationParameterSet = "Location";245 246 /// <summary>247 /// </summary>248 [Parameter]249 public SwitchParameter Wait250 {251 get252 {253 return _wait;254 }255 256 set257 {258 _wait = value;259 ValidateWait();260 }261 }262 263 /// <summary>264 /// </summary>265 [Parameter]266 public SwitchParameter AutoRemoveJob267 {268 get269 {270 return _autoRemoveJob;271 }272 273 set274 {275 _autoRemoveJob = value;276 }277 }278 279 /// <summary>280 /// </summary>281 [Parameter]282 public SwitchParameter WriteEvents283 {284 get285 {286 return _writeStateChangedEvents;287 }288 289 set290 {291 _writeStateChangedEvents = value;292 }293 }294 295 /// <summary>296 /// </summary>297 [Parameter]298 public SwitchParameter WriteJobInResults299 {300 get301 {302 return _outputJobFirst;303 }304 305 set306 {307 _outputJobFirst = value;308 }309 }310 311 private bool _autoRemoveJob;312 private bool _writeStateChangedEvents;313 private bool _wait;314 private bool _isStopping;315 private bool _isDisposed;316 private readonly ReaderWriterLockSlim _resultsReaderWriterLock = new ReaderWriterLockSlim();317 private readonly PowerShellTraceSource _tracer = PowerShellTraceSourceFactory.GetTraceSource();318 private readonly ManualResetEvent _writeExistingData = new ManualResetEvent(true);319 private readonly PSDataCollection<PSStreamObject> _results = new PSDataCollection<PSStreamObject>();320 private bool _holdingResultsRef;321 private readonly List<Job> _jobsBeingAggregated = new List<Job>();322 private readonly List<Guid> _jobsSpecifiedInParameters = new List<Guid>();323 private readonly object _syncObject = new object();324 private bool _outputJobFirst;325 private OutputProcessingState _outputProcessingNotification;326 private bool _processingOutput;327 328 private const string ClassNameTrace = "ReceiveJobCommand";329 330 #endregion Properties331 332 #region Overrides333 334 /// <summary>335 /// </summary>336 protected override void BeginProcessing()337 {338 ValidateAutoRemove();339 ValidateWriteJobInResults();340 ValidateWriteEvents();341 ValidateForce();342 }343 344 /// <summary>345 /// Retrieve the results for the specified computers or346 /// runspaces.347 /// </summary>348 protected override void ProcessRecord()349 {350 bool checkForRecurse = false;351 List<Job> jobsToWrite = new List<Job>();352 353 switch (ParameterSetName)354 {355 case SessionParameterSet:356 {357 foreach (Job job in _jobs)358 {359 PSRemotingJob remoteJob =360 job as PSRemotingJob;361 362 if (remoteJob == null)363 {364 string message = GetMessage(RemotingErrorIdStrings.RunspaceParamNotSupported);365 366 WriteError(new ErrorRecord(new ArgumentException(message),367 "RunspaceParameterNotSupported", ErrorCategory.InvalidArgument,368 job));369 370 continue;371 }372 373 // Runspace parameter is supported only on PSRemotingJob objects374 foreach (PSSession remoteRunspaceInfo in _remoteRunspaceInfos)375 {376 // get the required child jobs377 List<Job> childJobs = remoteJob.GetJobsForRunspace(remoteRunspaceInfo);378 jobsToWrite.AddRange(childJobs);379 // WriteResultsForJobsInCollection(childJobs, false);380 381 }382 }383 }384 385 break;386 387 case ComputerNameParameterSet:388 {389 foreach (Job job in _jobs)390 {391 // the job can either be a remoting job or another one392 PSRemotingJob remoteJob =393 job as PSRemotingJob;394 395 // ComputerName parameter can only be used with remoting jobs396 if (remoteJob == null)397 {398 string message = GetMessage(RemotingErrorIdStrings.ComputerNameParamNotSupported);399 400 WriteError(new ErrorRecord(new ArgumentException(message),401 "ComputerNameParameterNotSupported", ErrorCategory.InvalidArgument,402 job));403 404 continue;405 }406 407 string[] resolvedComputernames = null;408 ResolveComputerNames(_computerNames, out resolvedComputernames);409 410 foreach (string resolvedComputerName in resolvedComputernames)411 {412 // get the required child Job objects413 List<Job> childJobs = remoteJob.GetJobsForComputer(resolvedComputerName);414 jobsToWrite.AddRange(childJobs);415 // WriteResultsForJobsInCollection(childJobs, false);416 417 }418 }419 }420 421 break;422 423 case "Location":424 {425 if (_locations == null)426 {427 // WriteAll();428 jobsToWrite.AddRange(_jobs);429 checkForRecurse = true;430 }431 else432 {433 foreach (Job job in _jobs)434 {435 foreach (string location in _locations)436 {437 // get the required child Job objects438 List<Job> childJobs = job.GetJobsForLocation(location);439 jobsToWrite.AddRange(childJobs);440 // WriteResultsForJobsInCollection(childJobs, false);441 }442 }443 }444 }445 446 break;447 448 case ReceiveJobCommand.InstanceIdParameterSet:449 {450 List<Job> jobs = FindJobsMatchingByInstanceId(true, false, true, false);451 452 jobsToWrite.AddRange(jobs);453 checkForRecurse = true;454 // WriteResultsForJobsInCollection(jobs, true);455 }456 457 break;458 459 case ReceiveJobCommand.SessionIdParameterSet:460 {461 List<Job> jobs = FindJobsMatchingBySessionId(true, false, true, false);462 jobsToWrite.AddRange(jobs);463 checkForRecurse = true;464 // WriteResultsForJobsInCollection(jobs, true);465 }466 467 break;468 469 case ReceiveJobCommand.NameParameterSet:470 {471 List<Job> jobs = FindJobsMatchingByName(true, false, true, false);472 jobsToWrite.AddRange(jobs);473 checkForRecurse = true;474 // WriteResultsForJobsInCollection(jobs, true);475 }476 477 break;478 }479 480 // if block has been specified and the cmdlet has not been481 // stopped, we continue to write recursively, until there482 // is no more data to write483 if (_wait)484 {485 _writeExistingData.Reset();486 487 // if writejobresults is specified we will write only the top level jobs488 // this is because that is what the proxy requires. Anything else being489 // written is useless and will only add weight to the serialization490 491 WriteJobsIfRequired(jobsToWrite);492 493 // Make a note of the jobs specified by the user (does not include child jobs)494 // for the purpose of removal. Only the parent jobs should have remove called.495 foreach (var job in jobsToWrite)496 {497 _jobsSpecifiedInParameters.Add(job.InstanceId);498 }499 500 lock (_syncObject)501 {502 if (_isDisposed || _isStopping) return;503 504 // Check to see that we only AddRef once. ProcessRecord is called505 // once per job on the pipeline.506 if (!_holdingResultsRef)507 {508 _tracer.WriteMessage(ClassNameTrace, "ProcessRecord", Guid.Empty, (Job)null, "Adding Ref to results collection",509 null);510 _results.AddRef();511 _holdingResultsRef = true;512 }513 }514 515 _tracer.WriteMessage(ClassNameTrace, "ProcessRecord", Guid.Empty, (Job)null, "BEGIN Register for jobs");516 WriteResultsForJobsInCollection(jobsToWrite, checkForRecurse, true);517 _tracer.WriteMessage(ClassNameTrace, "ProcessRecord", Guid.Empty, (Job)null, "END Register for jobs");518 519 lock (_syncObject)520 {521 if (_jobsBeingAggregated.Count == 0 && _holdingResultsRef)522 {523 _tracer.WriteMessage(ClassNameTrace, "ProcessRecord", Guid.Empty, (Job)null,524 "Removing Ref to results collection", null);525 _results.DecrementRef();526 _holdingResultsRef = false;527 }528 }529 530 _tracer.WriteMessage(ClassNameTrace, "ProcessRecord", Guid.Empty, (Job)null, "BEGIN Write existing job data");531 WriteResultsForJobsInCollection(jobsToWrite, checkForRecurse, false);532 _tracer.WriteMessage(ClassNameTrace, "ProcessRecord", Guid.Empty, (Job)null, "END Write existing job data");533 _writeExistingData.Set();534 }535 else536 {537 WriteResultsForJobsInCollection(jobsToWrite, checkForRecurse, false);538 }539 }540 541 /// <summary>542 /// StopProcessing - when the command is stopped,543 /// unregister all the event handlers from the jobs544 /// and decrement reference for results.545 /// </summary>546 protected override void StopProcessing()547 {548 _tracer.WriteMessage(ClassNameTrace, "StopProcessing", Guid.Empty, (Job)null, "Entered Stop Processing",549 null);550 lock (_syncObject)551 {552 _isStopping = true;553 }554 555 _writeExistingData.Set();556 Job[] aggregatedJobs = new Job[_jobsBeingAggregated.Count];557 558 for (int i = 0; i < _jobsBeingAggregated.Count; i++)559 {560 aggregatedJobs[i] = _jobsBeingAggregated[i];561 }562 563 foreach (Job job in aggregatedJobs)564 {565 StopAggregateResultsFromJob(job);566 }567 568 _resultsReaderWriterLock.EnterWriteLock();569 try570 {571 _results.Complete();572 SetOutputProcessingState(false);573 }574 finally575 {576 _resultsReaderWriterLock.ExitWriteLock();577 }578 579 base.StopProcessing();580 _tracer.WriteMessage(ClassNameTrace, "StopProcessing", Guid.Empty, (Job)null, "Exiting Stop Processing",581 null);582 }583 584 /// <summary>585 /// If we are not stopping, continue writing output586 /// as and when they are available.587 /// </summary>588 protected override void EndProcessing()589 {590 try591 {592 if (_wait)593 {594 int totalCount = 0;595 foreach (PSStreamObject result in _results)596 {597 if (_isStopping) break;598 599 SetOutputProcessingState(true);600 result.WriteStreamObject(this, true, true);601 if (++totalCount == _results.Count)602 {603 SetOutputProcessingState(false);604 }605 }606 607 _eventArgsWritten.Clear();608 }609 else610 {611 int totalCount = 0;612 foreach (PSStreamObject result in _results)613 {614 if (_isStopping) break;615 616 SetOutputProcessingState(true);617 result.WriteStreamObject(this, false, true);618 if (++totalCount == _results.Count)619 {620 SetOutputProcessingState(false);621 }622 }623 }624 }625 finally626 {627 SetOutputProcessingState(false);628 }629 }630 631 /// <summary>632 /// </summary>633 public void Dispose()634 {635 Dispose(true);636 GC.SuppressFinalize(this);637 }638 639 /// <summary>640 /// </summary>641 /// <param name="disposing"></param>642 protected void Dispose(bool disposing)643 {644 if (disposing)645 {646 if (_isDisposed) return;647 lock (_syncObject)648 {649 if (_isDisposed) return;650 _isDisposed = true;651 }652 653 SetOutputProcessingState(false);654 655 if (_jobsBeingAggregated != null)656 {657 foreach (var job in _jobsBeingAggregated)658 {659 if (job.MonitorOutputProcessing)660 {661 job.RemoveMonitorOutputProcessing(_outputProcessingNotification);662 }663 664 if (job.UsesResultsCollection)665 {666 job.Results.DataAdded -= ResultsAdded;667 }668 else669 {670 job.Output.DataAdded -= Output_DataAdded;671 job.Error.DataAdded -= Error_DataAdded;672 job.Progress.DataAdded -= Progress_DataAdded;673 job.Verbose.DataAdded -= Verbose_DataAdded;674 job.Warning.DataAdded -= Warning_DataAdded;675 job.Debug.DataAdded -= Debug_DataAdded;676 job.Information.DataAdded -= Information_DataAdded;677 }678 679 job.StateChanged -= HandleJobStateChanged;680 }681 }682 683 _resultsReaderWriterLock.EnterWriteLock();684 try685 {686 _results.Complete();687 }688 finally689 {690 _resultsReaderWriterLock.ExitWriteLock();691 }692 693 _resultsReaderWriterLock.Dispose();694 _results.Clear();695 _results.Dispose();696 _writeExistingData.Set();697 _writeExistingData.Dispose();698 }699 }700 #endregion Overrides701 702 #region Private Methods703 704 private static void DoUnblockJob(Job job)705 {706 // we should not do anything for a parent job707 // the assumption is parent job states are708 // computed and so unblocking the child state709 // should be able to handle this710 if (job.ChildJobs.Count != 0) return;711 712 // we have a better way of handling blocked state logic713 // for remoting jobs, so use that if job is a remoting714 // job715 PSRemotingChildJob remotingChildJob = job as PSRemotingChildJob;716 if (remotingChildJob != null)717 {718 remotingChildJob.UnblockJob();719 }720 else721 {722 // for all other job types, simply set the job state723 // to running, the handling of the parent jobs state724 // should be taken care of by the job implementation725 job.SetJobState(JobState.Running, null);726 }727 }728 729 /// <summary>730 /// Write the results from this Job object. This does not write from the731 /// child jobs of this job object.732 /// </summary>733 /// <param name="job">Job object from which to write the results from734 /// </param>735 private void WriteJobResults(Job job)736 {737 if (job == null) return;738 739 // Q: Why do we need to unblock the job, before getting740 // the results741 // A: The job can get into a terminal state and we do742 // not want to set it to running at that point. Also, if743 // we do not explicitly signal that the job is unblocked744 // then the parent job cannot be unblocked. This is because745 // the parent job does not maintain a list of jobs which746 // are blocked but just simply a count (to keep things747 // light weight)748 749 // check if the state of the job is blocked, if so unblock it750 751 // Skip disconnected jobs that were in Blocked state before752 // the disconnect, since we cannot process host data until the753 // job is re-connected.754 if (job.JobStateInfo.State == JobState.Disconnected)755 {756 PSRemotingChildJob remotingChildJob = job as PSRemotingChildJob;757 if (remotingChildJob != null && remotingChildJob.DisconnectedAndBlocked)758 {759 return;760 }761 }762 763 // TODO: Fix Unblock() handling by Job2764 if (job.JobStateInfo.State == JobState.Blocked)765 {766 DoUnblockJob(job);767 }768 769 // for the jobs that PowerShell writes, there is a770 // results collection internally used. This collection771 // can be used to write results. For all other jobs772 // results need to be written from the other collections773 // available.774 // There is a bug in V2 that only remoting jobs work775 // with Receive-Job. This is being fixed776 777 if (job is not Job2 && job.UsesResultsCollection)778 {779 // extract results and handle them780 Collection<PSStreamObject> results = ReadAll<PSStreamObject>(job.Results);781 782 if (_wait)783 {784 foreach (var psStreamObject in results)785 {786 psStreamObject.WriteStreamObject(this, job.Results.SourceId);787 }788 }789 else790 {791 foreach (var psStreamObject in results)792 {793 psStreamObject.WriteStreamObject(this);794 }795 }796 }797 else798 {799 Collection<PSObject> output = ReadAll<PSObject>(job.Output);800 801 foreach (PSObject o in output)802 {803 if (o == null) continue;804 WriteObject(o);805 }806 807 Collection<ErrorRecord> errorRecords = ReadAll<ErrorRecord>(job.Error);808 809 foreach (ErrorRecord e in errorRecords)810 {811 if (e == null) continue;812 MshCommandRuntime mshCommandRuntime = CommandRuntime as MshCommandRuntime;813 if (mshCommandRuntime != null)814 {815 e.PreserveInvocationInfoOnce = true;816 mshCommandRuntime.WriteError(e, true);817 }818 }819 820 Collection<VerboseRecord> verboseRecords = ReadAll(job.Verbose);821 822 foreach (VerboseRecord v in verboseRecords)823 {824 if (v == null) continue;825 MshCommandRuntime mshCommandRuntime = CommandRuntime as MshCommandRuntime;826 mshCommandRuntime?.WriteVerbose(v, true);827 }828 829 Collection<DebugRecord> debugRecords = ReadAll(job.Debug);830 831 foreach (DebugRecord d in debugRecords)832 {833 if (d == null) continue;834 MshCommandRuntime mshCommandRuntime = CommandRuntime as MshCommandRuntime;835 mshCommandRuntime?.WriteDebug(d, true);836 }837 838 Collection<WarningRecord> warningRecords = ReadAll(job.Warning);839 840 foreach (WarningRecord w in warningRecords)841 {842 if (w == null) continue;843 MshCommandRuntime mshCommandRuntime = CommandRuntime as MshCommandRuntime;844 mshCommandRuntime?.WriteWarning(w, true);845 }846 847 Collection<ProgressRecord> progressRecords = ReadAll(job.Progress);848 849 foreach (ProgressRecord p in progressRecords)850 {851 if (p == null) continue;852 MshCommandRuntime mshCommandRuntime = CommandRuntime as MshCommandRuntime;853 mshCommandRuntime?.WriteProgress(p, true);854 }855 856 Collection<InformationRecord> informationRecords = ReadAll(job.Information);857 858 foreach (InformationRecord p in informationRecords)859 {860 if (p == null) continue;861 MshCommandRuntime mshCommandRuntime = CommandRuntime as MshCommandRuntime;862 mshCommandRuntime?.WriteInformation(p, true);863 }864 }865 866 if (job.JobStateInfo.State != JobState.Failed) return;867 868 WriteReasonError(job);869 }870 871 private void WriteReasonError(Job job)872 {873 // Write better error for the remoting case and generic error for the other case874 PSRemotingChildJob child = job as PSRemotingChildJob;875 if (child != null && child.FailureErrorRecord != null)876 {877 _results.Add(new PSStreamObject(PSStreamObjectType.Error, child.FailureErrorRecord, child.InstanceId));878 }879 else if (job.JobStateInfo.Reason != null)880 {881 Exception baseReason = job.JobStateInfo.Reason;882 Exception resultReason = baseReason;883 884 // If it was generated by a job that gave location information, unpack the885 // base exception.886 JobFailedException exceptionWithLocation = baseReason as JobFailedException;887 if (exceptionWithLocation != null)888 {889 resultReason = exceptionWithLocation.Reason;890 }891 892 ErrorRecord errorRecord = new ErrorRecord(resultReason, "JobStateFailed", ErrorCategory.InvalidResult, null);893 894 // If it was generated by a job that gave location information, set the895 // location information.896 if ((exceptionWithLocation != null) && (exceptionWithLocation.DisplayScriptPosition != null))897 {898 if (errorRecord.InvocationInfo == null)899 {900 errorRecord.SetInvocationInfo(new InvocationInfo(null, null));901 }902 903 errorRecord.InvocationInfo.DisplayScriptPosition = exceptionWithLocation.DisplayScriptPosition;904 }905 906 _results.Add(new PSStreamObject(PSStreamObjectType.Error, errorRecord, job.InstanceId));907 }908 }909 910 /// <summary>911 /// Returns all the results from supplied PSDataCollection.912 /// </summary>913 /// <param name="psDataCollection">Data collection to read from.</param>914 /// <returns>Collection with copy of data.</returns>915 private Collection<T> ReadAll<T>(PSDataCollection<T> psDataCollection)916 {917 if (_flush)918 {919 return psDataCollection.ReadAll();920 }921 922 T[] array = new T[psDataCollection.Count];923 psDataCollection.CopyTo(array, 0);924 Collection<T> collection = new Collection<T>();925 foreach (T t in array)926 {927 collection.Add(t);928 }929 930 return collection;931 }932 933 /// <summary>934 /// Write the results from this Job object. It also writes the935 /// results from its child objects recursively.936 /// </summary>937 /// <param name="duplicate">Hashtable used for duplicate detection.</param>938 /// <param name="job">Job whose results are written.</param>939 /// <param name="registerInsteadOfWrite"></param>940 private void WriteJobResultsRecursivelyHelper(Hashtable duplicate, Job job, bool registerInsteadOfWrite)941 {942 // Check if this object is already visited. If not, add it to the cache943 if (duplicate.ContainsKey(job))944 {945 return;946 }947 948 duplicate.Add(job, job);949 950 // Write the results of child jobs951 IList<Job> childJobs = job.ChildJobs;952 953 foreach (Job childjob in childJobs)954 {955 WriteJobResultsRecursivelyHelper(duplicate, childjob, registerInsteadOfWrite);956 }957 958 if (registerInsteadOfWrite)959 {960 // at any point there will be only one thread which will have961 // access to an entry corresponding to a job962 // this is because of the way the synchronization happens963 // with the pipeline thread and event handler thread using964 // _writeExistingData965 _eventArgsWritten[job.InstanceId] = false;966 // register the job for future updates967 AggregateResultsFromJob(job);968 }969 else970 {971 // Write the results of this job972 WriteJobResults(job);973 WriteJobStateInformationIfRequired(job);974 }975 }976 977 /// <summary>978 /// Writes the job objects if required by the cmdlet.979 /// </summary>980 /// <param name="jobsToWrite">Collection of jobs to write.</param>981 /// <remarks>this method is intended to be called only from982 /// ProcessRecord. When any changes are made ensure that this983 /// contract is not broken</remarks>984 private void WriteJobsIfRequired(IEnumerable<Job> jobsToWrite)985 {986 if (!_outputJobFirst) return;987 988 foreach (var job in jobsToWrite)989 {990 _tracer.WriteMessage("ReceiveJobCommand", "WriteJobsIfRequired", Guid.Empty, job, "Writing job object as output", null);991 WriteObject(job);992 }993 }994 995 /// <summary>996 /// </summary>997 /// <param name="job"></param>998 /// <remarks>this method should always be called before999 /// writeExistingData is set in ProcessRecord</remarks>1000 private void AggregateResultsFromJob(Job job)1001 {1002 if ((!Force && job.IsPersistentState(job.JobStateInfo.State)) || (Force && job.IsFinishedState(job.JobStateInfo.State))) return;1003 job.StateChanged += HandleJobStateChanged;1004 1005 // Check after the state changed event has been subscribed to avoid a race1006 // with the job state. StopAggregate is called from the state changed handler, so-1007 // this could cause the job to never be removed from _jobsBeingAggregated1008 // and therefore the _results ref to never be decremented.1009 if ((!Force && job.IsPersistentState(job.JobStateInfo.State)) || (Force && job.IsFinishedState(job.JobStateInfo.State)))1010 {1011 job.StateChanged -= HandleJobStateChanged;1012 return;1013 }1014 1015 _tracer.WriteMessage(ClassNameTrace, "AggregateResultsFromJob", Guid.Empty, job,1016 "BEGIN Adding job for aggregation", null);1017 1018 // at this point, we can be sure that any job added to this1019 // collection will have a state changed event to a finished state.1020 _jobsBeingAggregated.Add(job);1021 1022 // Tag the output collection so that the instance ID can be added to the output objects when streaming.1023 if (job.UsesResultsCollection)1024 {1025 job.Results.SourceId = job.InstanceId;1026 job.Results.DataAdded += ResultsAdded;1027 }1028 else1029 {1030 job.Output.SourceId = job.InstanceId;1031 job.Error.SourceId = job.InstanceId;1032 job.Progress.SourceId = job.InstanceId;1033 job.Verbose.SourceId = job.InstanceId;1034 job.Warning.SourceId = job.InstanceId;1035 job.Debug.SourceId = job.InstanceId;1036 job.Information.SourceId = job.InstanceId;1037 1038 job.Output.DataAdded += Output_DataAdded;1039 job.Error.DataAdded += Error_DataAdded;1040 job.Progress.DataAdded += Progress_DataAdded;1041 job.Verbose.DataAdded += Verbose_DataAdded;1042 job.Warning.DataAdded += Warning_DataAdded;1043 job.Debug.DataAdded += Debug_DataAdded;1044 job.Information.DataAdded += Information_DataAdded;1045 }1046 1047 if (job.MonitorOutputProcessing)1048 {1049 if (_outputProcessingNotification == null)1050 {1051 lock (_syncObject)1052 {1053 _outputProcessingNotification ??= new OutputProcessingState();1054 }1055 }1056 1057 job.SetMonitorOutputProcessing(_outputProcessingNotification);1058 }1059 1060 _tracer.WriteMessage(ClassNameTrace, "AggregateResultsFromJob", Guid.Empty, job,1061 "END Adding job for aggregation", null);1062 }1063 1064 private void ResultsAdded(object sender, DataAddedEventArgs e)1065 {1066 lock (_syncObject)1067 {1068 if (_isDisposed) return;1069 }1070 1071 _writeExistingData.WaitOne();1072 PSDataCollection<PSStreamObject> results = sender as PSDataCollection<PSStreamObject>;1073 1074 Dbg.Assert(results != null, "PSDataCollection is raising an inappropriate event");1075 PSStreamObject record = GetData(results, e.Index);1076 1077 if (record != null)1078 {1079 record.Id = results.SourceId;1080 _results.Add(record);1081 }1082 }1083 1084 private void HandleJobStateChanged(object sender, JobStateEventArgs e)1085 {1086 Job job = sender as Job;1087 Dbg.Assert(job != null, "Job state info cannot be raised with reference to job");1088 1089 // waiting for existing data to be written ensures two things1090 // 1. that no aggregation for a job is in progress1091 // 2. the state information is written in the correct order1092 // as per the contract1093 _tracer.WriteMessage(ClassNameTrace, "HandleJobStateChanged", Guid.Empty, job,1094 "BEGIN wait for write existing data", null);1095 if (e.JobStateInfo.State != JobState.Running)1096 _writeExistingData.WaitOne();1097 1098 _tracer.WriteMessage(ClassNameTrace, "HandleJobStateChanged", Guid.Empty, job,1099 "END wait for write existing data", null);1100 1101 lock (_syncObject)1102 {1103 if (!_jobsBeingAggregated.Contains(job))1104 {1105 _tracer.WriteMessage(ClassNameTrace, "HandleJobStateChanged", Guid.Empty, job,1106 "Returning because job is not in _jobsBeingAggregated", null);1107 return;1108 }1109 }1110 1111 if (e.JobStateInfo.State == JobState.Blocked)1112 {1113 DoUnblockJob(job);1114 }1115 1116 // Stop wait if:1117 // Force is specified and Job is in a Finished state (Completed, Failed, Stopped)1118 // OR1119 // Force is not specified and Job is in a persistent state (Suspended or1120 // Disconnected as well as above)1121 // (logic inverted for return)1122 if (!(!Force && job.IsPersistentState(e.JobStateInfo.State)) && !(Force && job.IsFinishedState(e.JobStateInfo.State)))1123 {1124 _tracer.WriteMessage(ClassNameTrace, "HandleJobStateChanged", Guid.Empty, job,1125 "Returning because job state does not meet wait requirements (continue aggregating)");1126 return;1127 }1128 1129 // Write an error record with JobStateFailed ID if there is a JobStateInfo.Reason.1130 WriteReasonError(job);1131 WriteJobStateInformationIfRequired(job, e);1132 StopAggregateResultsFromJob(job);1133 }1134 1135 private void Progress_DataAdded(object sender, DataAddedEventArgs e)1136 {1137 lock (_syncObject)1138 {1139 if (_isDisposed) return;1140 }1141 1142 _writeExistingData.WaitOne();1143 _resultsReaderWriterLock.EnterReadLock();1144 try1145 {1146 if (!_results.IsOpen) return;1147 PSDataCollection<ProgressRecord> progressRecords = sender as PSDataCollection<ProgressRecord>;1148 Diagnostics.Assert(progressRecords != null, "PSDataCollection is raising an inappropriate event");1149 1150 ProgressRecord record = GetData(progressRecords, e.Index);1151 if (record != null)1152 {1153 _results.Add(new PSStreamObject(PSStreamObjectType.Progress, record, progressRecords.SourceId));1154 }1155 }1156 finally1157 {1158 _resultsReaderWriterLock.ExitReadLock();1159 }1160 }1161 1162 private void Error_DataAdded(object sender, DataAddedEventArgs e)1163 {1164 lock (_syncObject)1165 {1166 if (_isDisposed)1167 {1168 return;1169 }1170 }1171 1172 _writeExistingData.WaitOne();1173 _resultsReaderWriterLock.EnterReadLock();1174 try1175 {1176 if (!_results.IsOpen)1177 {1178 return;1179 }1180 1181 PSDataCollection<ErrorRecord> errorRecords = sender as PSDataCollection<ErrorRecord>;1182 Diagnostics.Assert(errorRecords != null, "PSDataCollection is raising an inappropriate event");1183 ErrorRecord errorRecord = GetData(errorRecords, e.Index);1184 if (errorRecord != null)1185 {1186 // error records are already tagged, skip tagging1187 _results.Add(new PSStreamObject(PSStreamObjectType.Error, errorRecord, Guid.Empty));1188 }1189 }1190 finally1191 {1192 _resultsReaderWriterLock.ExitReadLock();1193 }1194 }1195 1196 private void Debug_DataAdded(object sender, DataAddedEventArgs e)1197 {1198 lock (_syncObject)1199 {1200 if (_isDisposed) return;