Implement ForEach-Object -Parallel runspace reuse (#12122)

* Implement foreach parallel runspace reuse

* Change runspace dispose

* Refactor runspace reset check

* Fix race condition.

* Fix CodFactor issues

* Implement -UseNewRunspace parameter switch, add tests

* Fix Codacy error
This commit is contained in:
Paul Higinbotham
2020-03-23 14:00:43 -07:00
committed by GitHub
parent 7a8094fd31
commit cce214e884
3 changed files with 185 additions and 29 deletions
@@ -260,6 +260,14 @@ namespace Microsoft.PowerShell.Commands
[Parameter(ParameterSetName = ForEachObjectCommand.ParallelParameterSet)]
public SwitchParameter AsJob { get; set; }
/// <summary>
/// Gets or sets a flag so that a new runspace object is created for each loop iteration, instead of reusing objects
/// from the runspace pool.
/// By default, runspaces are reused from a runspace pool.
/// </summary>
[Parameter(ParameterSetName = ForEachObjectCommand.ParallelParameterSet)]
public SwitchParameter UseNewRunspace { get; set; }
#endregion
#region Overrides
@@ -437,7 +445,8 @@ namespace Microsoft.PowerShell.Commands
_taskJob = new PSTaskJob(
Parallel.ToString(),
ThrottleLimit);
ThrottleLimit,
UseNewRunspace);
return;
}
@@ -445,7 +454,7 @@ namespace Microsoft.PowerShell.Commands
// Set up for synchronous processing and data streaming.
_taskCollection = new PSDataCollection<System.Management.Automation.PSTasks.PSTask>();
_taskDataStreamWriter = new PSTaskDataStreamWriter(this);
_taskPool = new PSTaskPool(ThrottleLimit);
_taskPool = new PSTaskPool(ThrottleLimit, UseNewRunspace);
_taskPool.PoolComplete += (sender, args) =>
{
_taskDataStreamWriter.Close();
@@ -1,6 +1,7 @@
// Copyright (c) Microsoft Corporation. All rights reserved.
// Licensed under the MIT License.
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Globalization;
using System.Management.Automation.Host;
@@ -320,7 +321,7 @@ namespace System.Management.Automation.PSTasks
protected PowerShell _powershell;
protected PSDataCollection<PSObject> _output;
private const string RunspaceName = "PSTask";
public const string RunspaceName = "PSTask";
private static int s_taskId;
@@ -364,6 +365,11 @@ namespace System.Management.Automation.PSTasks
/// </summary>
public int Id { get => _id; }
/// <summary>
/// Gets Task Runspace.
/// </summary>
public Runspace Runspace { get => _runspace; }
#endregion
#region Constructor
@@ -410,7 +416,6 @@ namespace System.Management.Automation.PSTasks
/// </summary>
public void Dispose()
{
_runspace.Dispose();
_powershell.Dispose();
_output.Dispose();
}
@@ -422,7 +427,8 @@ namespace System.Management.Automation.PSTasks
/// <summary>
/// Start task.
/// </summary>
public void Start()
/// <param name="runspace">Runspace used to run task.</param>
public void Start(Runspace runspace)
{
if (_powershell != null)
{
@@ -430,13 +436,8 @@ namespace System.Management.Automation.PSTasks
return;
}
// Create and open Runspace for this task to run in
var iss = InitialSessionState.CreateDefault2();
iss.LanguageMode = (SystemPolicy.GetSystemLockdownPolicy() == SystemEnforcementMode.Enforce)
? PSLanguageMode.ConstrainedLanguage : PSLanguageMode.FullLanguage;
_runspace = RunspaceFactory.CreateRunspace(iss);
_runspace.Name = string.Format(CultureInfo.InvariantCulture, "{0}:{1}", RunspaceName, s_taskId);
_runspace.Open();
Dbg.Assert(runspace != null, "Task runspace cannot be null.");
_runspace = runspace;
// If available, set current working directory on the runspace.
// Temporarily set the newly created runspace as the thread default runspace for any needed module loading.
@@ -445,8 +446,8 @@ namespace System.Management.Automation.PSTasks
var oldDefaultRunspace = Runspace.DefaultRunspace;
try
{
Runspace.DefaultRunspace = _runspace;
_runspace.ExecutionContext.SessionState.Internal.SetLocation(_currentLocationPath);
Runspace.DefaultRunspace = runspace;
runspace.ExecutionContext.SessionState.Internal.SetLocation(_currentLocationPath);
}
finally
{
@@ -456,7 +457,7 @@ namespace System.Management.Automation.PSTasks
// Create the PowerShell command pipeline for the provided script block
// The script will run on the provided Runspace in a new thread by default
_powershell = PowerShell.Create(_runspace);
_powershell = PowerShell.Create(runspace);
// Initialize PowerShell object data streams and event handlers
_output = new PSDataCollection<PSObject>();
@@ -632,8 +633,13 @@ namespace System.Management.Automation.PSTasks
private readonly ManualResetEvent _stopAll;
private readonly object _syncObject;
private readonly Dictionary<int, PSTaskBase> _taskPool;
private readonly ConcurrentQueue<Runspace> _runspacePool;
private readonly ConcurrentDictionary<int, Runspace> _activeRunspaces;
private readonly WaitHandle[] _waitHandles;
private readonly bool _useRunspacePool;
private bool _isOpen;
private bool _stopping;
private int _createdRunspaceCount;
private const int AddAvailable = 0;
private const int Stop = 1;
@@ -648,9 +654,13 @@ namespace System.Management.Automation.PSTasks
/// Initializes a new instance of the <see cref="PSTaskPool"/> class.
/// </summary>
/// <param name="size">Total number of allowed running objects in pool at one time.</param>
public PSTaskPool(int size)
/// <param name="useNewRunspace">When true, a new runspace object is created for the task instead of reusing one from the pool.</param>
public PSTaskPool(
int size,
bool useNewRunspace)
{
_sizeLimit = size;
_useRunspacePool = !useNewRunspace;
_isOpen = true;
_syncObject = new object();
_addAvailable = new ManualResetEvent(true);
@@ -661,6 +671,11 @@ namespace System.Management.Automation.PSTasks
_stopAll, // index 1
};
_taskPool = new Dictionary<int, PSTaskBase>(size);
_activeRunspaces = new ConcurrentDictionary<int, Runspace>();
if (_useRunspacePool)
{
_runspacePool = new ConcurrentQueue<Runspace>();
}
}
#endregion
@@ -684,6 +699,14 @@ namespace System.Management.Automation.PSTasks
get => _isOpen;
}
/// <summary>
/// Gets a value of the count of total runspaces allocated.
/// </summary>
public int AllocatedRunspaceCount
{
get => _createdRunspaceCount;
}
#endregion
#region IDisposable
@@ -695,6 +718,21 @@ namespace System.Management.Automation.PSTasks
{
_addAvailable.Dispose();
_stopAll.Dispose();
DisposeRunspaces();
}
/// <summary>
/// Dispose runspaces.
/// </summary>
internal void DisposeRunspaces()
{
foreach (var item in _activeRunspaces)
{
item.Value.Dispose();
}
_activeRunspaces.Clear();
}
#endregion
@@ -721,6 +759,7 @@ namespace System.Management.Automation.PSTasks
switch (index)
{
case AddAvailable:
var runspace = GetRunspace(task.Id);
task.StateChanged += HandleTaskStateChangedDelegate;
lock (_syncObject)
{
@@ -735,7 +774,7 @@ namespace System.Management.Automation.PSTasks
_addAvailable.Reset();
}
task.Start();
task.Start(runspace);
}
return true;
@@ -763,18 +802,28 @@ namespace System.Management.Automation.PSTasks
/// </summary>
public void StopAll()
{
_stopping = true;
// Accept no more input
Close();
_stopAll.Set();
// Stop all running tasks
PSTaskBase[] tasksToStop;
lock (_syncObject)
{
foreach (var task in _taskPool.Values)
{
task.Dispose();
}
tasksToStop = new PSTaskBase[_taskPool.Values.Count];
_taskPool.Values.CopyTo(tasksToStop, 0);
}
foreach (var task in tasksToStop)
{
task.Dispose();
}
// Dispose all active runspaces
DisposeRunspaces();
_stopping = false;
}
/// <summary>
@@ -803,6 +852,7 @@ namespace System.Management.Automation.PSTasks
case PSInvocationState.Completed:
case PSInvocationState.Stopped:
case PSInvocationState.Failed:
ReturnRunspace(task);
lock (_syncObject)
{
_taskPool.Remove(task.Id);
@@ -813,7 +863,12 @@ namespace System.Management.Automation.PSTasks
}
task.StateChanged -= HandleTaskStateChangedDelegate;
task.Dispose();
if (!_stopping)
{
// StopAll disposes tasks.
task.Dispose();
}
CheckForComplete();
break;
}
@@ -842,6 +897,64 @@ namespace System.Management.Automation.PSTasks
}
}
private Runspace GetRunspace(int taskId)
{
var runspaceName = string.Format(CultureInfo.InvariantCulture, "{0}:{1}", PSTask.RunspaceName, taskId);
if (_useRunspacePool && _runspacePool.TryDequeue(out Runspace runspace))
{
if (runspace.RunspaceStateInfo.State == RunspaceState.Opened &&
runspace.RunspaceAvailability == RunspaceAvailability.Available)
{
try
{
runspace.ResetRunspaceState();
runspace.Name = runspaceName;
return runspace;
}
catch
{
// If the runspace cannot be reset for any reason, remove it.
}
}
RemoveActiveRunspace(runspace);
}
// Create and initialize a new Runspace
var iss = InitialSessionState.CreateDefault2();
iss.LanguageMode = (SystemPolicy.GetSystemLockdownPolicy() == SystemEnforcementMode.Enforce)
? PSLanguageMode.ConstrainedLanguage : PSLanguageMode.FullLanguage;
runspace = RunspaceFactory.CreateRunspace(iss);
runspace.Name = runspaceName;
_activeRunspaces.TryAdd(runspace.Id, runspace);
runspace.Open();
_createdRunspaceCount++;
return runspace;
}
private void ReturnRunspace(PSTaskBase task)
{
var runspace = task.Runspace;
Dbg.Assert(runspace != null, "Task runspace cannot be null.");
if (_useRunspacePool &&
runspace.RunspaceStateInfo.State == RunspaceState.Opened &&
runspace.RunspaceAvailability == RunspaceAvailability.Available)
{
_runspacePool.Enqueue(runspace);
return;
}
RemoveActiveRunspace(runspace);
}
private void RemoveActiveRunspace(Runspace runspace)
{
runspace.Dispose();
_activeRunspaces.TryRemove(runspace.Id, out Runspace _);
}
#endregion
}
@@ -852,7 +965,7 @@ namespace System.Management.Automation.PSTasks
/// <summary>
/// Job for running ForEach-Object parallel task child jobs asynchronously.
/// </summary>
internal sealed class PSTaskJob : Job
public sealed class PSTaskJob : Job
{
#region Members
@@ -862,6 +975,18 @@ namespace System.Management.Automation.PSTasks
#endregion
#region Properties
/// <summary>
/// Gets a value of the count of total runspaces allocated.
/// </summary>
public int AllocatedRunspaceCount
{
get => _taskPool.AllocatedRunspaceCount;
}
#endregion
#region Constructor
private PSTaskJob() { }
@@ -871,11 +996,13 @@ namespace System.Management.Automation.PSTasks
/// </summary>
/// <param name="command">Job command text.</param>
/// <param name="throttleLimit">Pool size limit for task job.</param>
public PSTaskJob(
/// <param name="useNewRunspace">When true, a new runspace object is created for the task instead of reusing one from the pool.</param>
internal PSTaskJob(
string command,
int throttleLimit) : base(command, string.Empty)
int throttleLimit,
bool useNewRunspace) : base(command, string.Empty)
{
_taskPool = new PSTaskPool(throttleLimit);
_taskPool = new PSTaskPool(throttleLimit, useNewRunspace);
_isOpen = true;
PSJobTypeName = nameof(PSTaskJob);
@@ -949,14 +1076,14 @@ namespace System.Management.Automation.PSTasks
#endregion
#region Public Methods
#region Internal Methods
/// <summary>
/// Add a child job to the collection.
/// </summary>
/// <param name="childJob">Child job to add.</param>
/// <returns>True when child job is successfully added.</returns>
public bool AddJob(PSTaskChildJob childJob)
internal bool AddJob(PSTaskChildJob childJob)
{
if (!_isOpen)
{
@@ -971,7 +1098,7 @@ namespace System.Management.Automation.PSTasks
/// Closes this parent job to adding more child jobs and starts
/// the child jobs running with the provided throttle limit.
/// </summary>
public void Start()
internal void Start()
{
_isOpen = false;
SetJobState(JobState.Running);
@@ -1017,6 +1144,9 @@ namespace System.Management.Automation.PSTasks
}
SetJobState(finalState);
// Release job task pool runspace resources.
(sender as PSTaskPool).DisposeRunspaces();
}
catch (ObjectDisposedException)
{ }
@@ -314,6 +314,23 @@ Describe 'ForEach-Object -Parallel -AsJob Basic Tests' -Tags 'CI' {
}
}
Describe 'ForEach-Object -Parallel runspace pool tests' -Tags 'CI' {
It "Verifies job allocated runspace count is limited to pool size" {
$job = 1..4 | ForEach-Object -Parallel { Start-Sleep 1 } -AsJob -ThrottleLimit 2 | Wait-Job
$job.AllocatedRunspaceCount | Should -BeExactly 2
$job | Remove-Job
}
It "Verifies job with -UseNewRunspace switch allocates one runspace per iteration" {
$job = 1..10 | ForEach-Object -Parallel { $_ } -AsJob -ThrottleLimit 2 -UseNewRunspace | Wait-Job
$job.AllocatedRunspaceCount | Should -BeExactly 10
$job | Remove-Job
}
}
Describe 'ForEach-Object -Parallel Functional Tests' -Tags 'Feature' {
It 'Verifies job queuing and throttle limit' {