Team Ai
Datasetpublic

MegaBites-AI/Windows-powershell

sourceHugging Facemitupdated 6mo agoView on Hugging Face
0likes372downloads
ReceiveJob.cs1614 linesDownload Raw Back to commands
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;

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