Refactor ProcessScheduler with Execute method

This commit is contained in:
jonawals
2021-05-21 20:13:10 +01:00
parent 444c0839db
commit 597f0905c5
12 changed files with 272 additions and 209 deletions
@@ -50,32 +50,29 @@ namespace TestImpact
{
public:
//! Constructs the job runner with the specified parameters to constrain job runs.
//! @param maxConcurrentProcesses he maximum number of concurrent jobs in-flight.
JobRunner(size_t maxConcurrentProcesses);
//! Executes the specified jobs and returns the products of their labor.
//! @param jobs The arguments (and other pertinent information) required for each job to be run.
//! @param stdOutRouting The standard output routing to be specified for all jobs.
//! @param stdErrRouting The standard error routing to be specified for all jobs.
//! @param maxConcurrentProcesses he maximum number of concurrent jobs in-flight.
//! @param processTimeout The maximum duration a job may be in-flight before being forcefully terminated (nullopt if no timeout).
//! @param scheduleTimeout The maximum duration the scheduler may run before forcefully terminating all in-flight jobs (nullopt if
//! no timeout).
JobRunner(
//! @param jobTimeout The maximum duration a job may be in-flight before being forcefully terminated (nullopt if no timeout).
//! @param runnerTimeout The maximum duration the scheduler may run before forcefully terminating all in-flight jobs (nullopt if no timeout).
//! @param payloadMapProducer The client callback to be called when all jobs have finished to transform the work produced by each job into the desired output.
//! @param jobCallback The client callback to be called when each job changes state.
//! @return The result of the run sequence and the jobs with their associated payloads.
AZStd::pair<ProcessSchedulerResult, AZStd::vector<typename JobT>> Execute(
const AZStd::vector<typename JobT::Info>& jobs,
PayloadMapProducer<JobT> payloadMapProducer,
StdOutputRouting stdOutRouting,
StdErrorRouting stdErrRouting,
size_t maxConcurrentProcesses,
AZStd::optional<AZStd::chrono::milliseconds> processTimeout,
AZStd::optional<AZStd::chrono::milliseconds> scheduleTimeout);
//! Executes the specified jobs and returns the products of their labor.
//! @note: the job and payload callbacks are specified here rather than in the constructor to allow clients to use capturing lambdas
//! should they desire to.
//! @param jobs The arguments (and other pertinent information) required for each job to be run.
//! @param jobCallback The client callback to be called when each job changes state.
//! @param payloadMapProducer The client callback to be called when all jobs have finished to transform the work produced by each
//! job into the desired output.
AZStd::vector<typename JobT> Execute(
const AZStd::vector<typename JobT::Info>& jobs, JobCallback<typename JobT> jobCallback,
PayloadMapProducer<JobT> payloadMapProducer);
AZStd::optional<AZStd::chrono::milliseconds> jobTimeout,
AZStd::optional<AZStd::chrono::milliseconds> runnerTimeout,
JobCallback<typename JobT> jobCallback);
private:
size_t m_maxConcurrentProcesses = 0; //!< Maximum number of concurrent jobs being executed at a given time.
ProcessScheduler m_processScheduler;
StdOutputRouting m_stdOutRouting; //!< Standard output routing from each job process to job runner.
StdErrorRouting m_stdErrRouting; //!< Standard error routing from each job process to job runner
AZStd::optional<AZStd::chrono::milliseconds> m_jobTimeout; //!< Maximum time a job can run for before being forcefully terminated.
@@ -83,25 +80,20 @@ namespace TestImpact
};
template<typename JobT>
JobRunner<JobT>::JobRunner(
StdOutputRouting stdOutRouting,
StdErrorRouting stdErrRouting,
size_t maxConcurrentProcesses,
AZStd::optional<AZStd::chrono::milliseconds> jobTimeout,
AZStd::optional<AZStd::chrono::milliseconds> runnerTimeout)
: m_maxConcurrentProcesses(maxConcurrentProcesses)
, m_stdOutRouting(stdOutRouting)
, m_stdErrRouting(stdErrRouting)
, m_jobTimeout(jobTimeout)
, m_runnerTimeout(runnerTimeout)
JobRunner<JobT>::JobRunner(size_t maxConcurrentProcesses)
: m_processScheduler(maxConcurrentProcesses)
{
}
template<typename JobT>
AZStd::vector<JobT> JobRunner<JobT>::Execute(
AZStd::pair<ProcessSchedulerResult, AZStd::vector<typename JobT>> JobRunner<JobT>::Execute(
const AZStd::vector<typename JobT::Info>& jobInfos,
JobCallback<JobT> jobCallback,
PayloadMapProducer<JobT> payloadMapProducer)
PayloadMapProducer<JobT> payloadMapProducer,
StdOutputRouting stdOutRouting,
StdErrorRouting stdErrRouting,
AZStd::optional<AZStd::chrono::milliseconds> jobTimeout,
AZStd::optional<AZStd::chrono::milliseconds> runnerTimeout,
JobCallback<typename JobT> jobCallback)
{
AZStd::vector<ProcessInfo> processes;
AZStd::unordered_map<JobT::Info::IdType, AZStd::pair<JobMeta, const typename JobT::Info*>> metas;
@@ -115,11 +107,11 @@ namespace TestImpact
const auto* jobInfo = &jobInfos[jobIndex];
const auto jobId = jobInfo->GetId().m_value;
metas.emplace(jobId, AZStd::pair<JobMeta, const typename JobT::Info*>{JobMeta{}, jobInfo});
processes.emplace_back(jobId, m_stdOutRouting, m_stdErrRouting, jobInfo->GetCommand().m_args);
processes.emplace_back(jobId, stdOutRouting, stdErrRouting, jobInfo->GetCommand().m_args);
}
// Wrapper around low-level process launch callback to gather job meta-data and present a simplified callback interface to the client
const auto processLaunchCallback = [&jobCallback, &jobInfos, &metas](
const ProcessLaunchCallback processLaunchCallback = [&jobCallback, &jobInfos, &metas](
TestImpact::ProcessId pid,
TestImpact::LaunchResult launchResult,
AZStd::chrono::high_resolution_clock::time_point createTime)
@@ -138,7 +130,7 @@ namespace TestImpact
};
// Wrapper around low-level process exit callback to gather job meta-data and present a simplified callback interface to the client
const auto processExitCallback = [&jobCallback, &jobInfos, &metas](
const ProcessExitCallback processExitCallback = [&jobCallback, &jobInfos, &metas](
TestImpact::ProcessId pid,
TestImpact::ExitCondition exitCondition,
TestImpact::ReturnCode returnCode,
@@ -169,14 +161,12 @@ namespace TestImpact
};
// Schedule all jobs for execution
ProcessScheduler scheduler(
const auto result = m_processScheduler.Execute(
processes,
jobTimeout,
runnerTimeout,
processLaunchCallback,
processExitCallback,
m_maxConcurrentProcesses,
m_jobTimeout,
m_runnerTimeout
);
processExitCallback);
// Hand off the jobs to the client for payload generation
auto payloadMap = payloadMapProducer(metas);
@@ -188,6 +178,6 @@ namespace TestImpact
jobs.emplace_back(JobT(jobInfo, AZStd::move(metas.at(jobId).first), AZStd::move(payloadMap[jobId])));
}
return jobs;
return { result, jobs };
}
} // namespace TestImpact
@@ -18,7 +18,7 @@
namespace TestImpact
{
struct ProcessScheduler::ProcessInFlight
struct ProcessInFlight
{
AZStd::unique_ptr<Process> m_process;
AZStd::optional<AZStd::chrono::high_resolution_clock::time_point> m_startTime;
@@ -26,29 +26,64 @@ namespace TestImpact
AZStd::string m_stdError;
};
ProcessScheduler::ProcessScheduler(
const AZStd::vector<ProcessInfo>& processes,
const ProcessLaunchCallback& processLaunchCallback,
const ProcessExitCallback& processExitCallback,
class ProcessScheduler::ExecutionState
{
public:
ExecutionState(
size_t maxConcurrentProcesses,
AZStd::optional<AZStd::chrono::milliseconds> processTimeout,
AZStd::optional<AZStd::chrono::milliseconds> scheduleTimeout,
ProcessLaunchCallback& processLaunchCallback,
ProcessExitCallback& processExitCallback);
~ExecutionState();
ProcessSchedulerResult MonitorProcesses(const AZStd::vector<ProcessInfo>& processes);
void TerminateAllProcesses(ExitCondition exitStatus);
private:
ProcessCallbackResult PopAndLaunch(ProcessInFlight& processInFlight);
StdContent ConsumeProcessStdContent(ProcessInFlight& processInFlight);
void AccumulateProcessStdContent(ProcessInFlight& processInFlight);
size_t m_maxConcurrentProcesses = 0;
ProcessLaunchCallback m_processLaunchCallback;
ProcessExitCallback m_processExitCallback;
AZStd::optional<AZStd::chrono::milliseconds> m_processTimeout;
AZStd::optional<AZStd::chrono::milliseconds> m_scheduleTimeout;
AZStd::chrono::high_resolution_clock::time_point m_startTime;
AZStd::vector<ProcessInFlight> m_processPool;
AZStd::queue<ProcessInfo> m_processQueue;
};
ProcessScheduler::ExecutionState::ExecutionState(
size_t maxConcurrentProcesses,
AZStd::optional<AZStd::chrono::milliseconds> processTimeout,
AZStd::optional<AZStd::chrono::milliseconds> scheduleTimeout)
: m_processCreateCallback(processLaunchCallback)
AZStd::optional<AZStd::chrono::milliseconds> scheduleTimeout,
ProcessLaunchCallback& processLaunchCallback,
ProcessExitCallback& processExitCallback)
: m_maxConcurrentProcesses(maxConcurrentProcesses)
, m_processLaunchCallback(processLaunchCallback)
, m_processExitCallback(processExitCallback)
, m_processTimeout(processTimeout)
, m_scheduleTimeout(scheduleTimeout)
, m_startTime(AZStd::chrono::high_resolution_clock::now())
{
AZ_TestImpact_Eval(maxConcurrentProcesses != 0, ProcessException, "Max Number of concurrent processes in flight cannot be 0");
AZ_TestImpact_Eval(!processes.empty(), ProcessException, "Number of processes to launch cannot be 0");
AZ_TestImpact_Eval(
!m_processTimeout.has_value() || m_processTimeout->count() > 0, ProcessException,
"Process timeout must be empty or non-zero value");
AZ_TestImpact_Eval(
!m_scheduleTimeout.has_value() || m_scheduleTimeout->count() > 0, ProcessException,
"Scheduler timeout must be empty or non-zero value");
}
const size_t numConcurrentProcesses = AZStd::min(processes.size(), maxConcurrentProcesses);
ProcessScheduler::ExecutionState::~ExecutionState()
{
TerminateAllProcesses(ExitCondition::Terminated);
}
ProcessSchedulerResult ProcessScheduler::ExecutionState::MonitorProcesses(const AZStd::vector<ProcessInfo>& processes)
{
AZ_TestImpact_Eval(!processes.empty(), ProcessException, "Number of processes to launch cannot be 0");
m_startTime = AZStd::chrono::high_resolution_clock::now();
const size_t numConcurrentProcesses = AZStd::min(processes.size(), m_maxConcurrentProcesses);
m_processPool.resize(numConcurrentProcesses);
for (const auto& process : processes)
@@ -61,20 +96,10 @@ namespace TestImpact
if (PopAndLaunch(process) == ProcessCallbackResult::Abort)
{
TerminateAllProcesses(ExitCondition::Terminated);
return;
return ProcessSchedulerResult::Graceful;
}
}
MonitorProcesses();
}
ProcessScheduler::~ProcessScheduler()
{
TerminateAllProcesses(ExitCondition::Terminated);
}
void ProcessScheduler::MonitorProcesses()
{
while (true)
{
// Check to see whether or not the scheduling has exceeded its specified runtime
@@ -86,7 +111,7 @@ namespace TestImpact
{
// Runtime exceeded, terminate all proccesses and schedule no further
TerminateAllProcesses(ExitCondition::Timeout);
return;
return ProcessSchedulerResult::Timeout;
}
}
@@ -98,7 +123,7 @@ namespace TestImpact
{
if (processInFlight.m_process)
{
// Process is alive (note: not necessarilly currently running)
// Process is alive (note: not necessarily currently running)
AccumulateProcessStdContent(processInFlight);
const ProcessId processId = processInFlight.m_process->GetProcessInfo().GetId();
@@ -119,7 +144,7 @@ namespace TestImpact
{
// Client chose to abort the scheduler
TerminateAllProcesses(ExitCondition::Terminated);
return;
return ProcessSchedulerResult::Graceful;
}
else if (!m_processQueue.empty())
{
@@ -128,7 +153,7 @@ namespace TestImpact
{
// Client chose to abort the scheduler
TerminateAllProcesses(ExitCondition::Terminated);
return;
return ProcessSchedulerResult::Graceful;
}
else
{
@@ -157,9 +182,9 @@ namespace TestImpact
ConsumeProcessStdContent(processInFlight),
exitTime))
{
// Flight time exceeded, terminate this process
// Client chose to abort the scheduler
TerminateAllProcesses(ExitCondition::Terminated);
return;
return ProcessSchedulerResult::Graceful;
}
}
@@ -176,7 +201,7 @@ namespace TestImpact
{
// Client chose to abort the scheduler
TerminateAllProcesses(ExitCondition::Terminated);
return;
return ProcessSchedulerResult::Graceful;
}
else
{
@@ -192,9 +217,11 @@ namespace TestImpact
break;
}
}
return ProcessSchedulerResult::Graceful;
}
ProcessCallbackResult ProcessScheduler::PopAndLaunch(ProcessInFlight& processInFlight)
ProcessCallbackResult ProcessScheduler::ExecutionState::PopAndLaunch(ProcessInFlight& processInFlight)
{
auto processInfo = m_processQueue.front();
m_processQueue.pop();
@@ -212,17 +239,17 @@ namespace TestImpact
createResult = LaunchResult::Failure;
}
return m_processCreateCallback(processInfo.GetId(), createResult, createTime);
return m_processLaunchCallback(processInfo.GetId(), createResult, createTime);
}
void ProcessScheduler::AccumulateProcessStdContent(ProcessInFlight& processInFlight)
void ProcessScheduler::ExecutionState::AccumulateProcessStdContent(ProcessInFlight& processInFlight)
{
// Accumulate the stdout/stderr so we don't deadlock with the process waiting for the pipe to empty before finishing
processInFlight.m_stdOutput += processInFlight.m_process->ConsumeStdOut().value_or("");
processInFlight.m_stdError += processInFlight.m_process->ConsumeStdErr().value_or("");
}
StdContent ProcessScheduler::ConsumeProcessStdContent(ProcessInFlight& processInFlight)
StdContent ProcessScheduler::ExecutionState::ConsumeProcessStdContent(ProcessInFlight& processInFlight)
{
return
{
@@ -235,7 +262,7 @@ namespace TestImpact
};
}
void ProcessScheduler::TerminateAllProcesses(ExitCondition exitStatus)
void ProcessScheduler::ExecutionState::TerminateAllProcesses(ExitCondition exitStatus)
{
bool isCallingBackToClient = true;
const ReturnCode returnCode = static_cast<ReturnCode>(exitStatus);
@@ -267,4 +294,27 @@ namespace TestImpact
}
}
}
ProcessScheduler::ProcessScheduler(size_t maxConcurrentProcesses)
: m_maxConcurrentProcesses(maxConcurrentProcesses)
{
AZ_TestImpact_Eval(maxConcurrentProcesses != 0, ProcessException, "Max Number of concurrent processes in flight cannot be 0");
}
ProcessScheduler::~ProcessScheduler() = default;
ProcessSchedulerResult ProcessScheduler::Execute(
const AZStd::vector<ProcessInfo>& processes,
AZStd::optional<AZStd::chrono::milliseconds> processTimeout,
AZStd::optional<AZStd::chrono::milliseconds> scheduleTimeout,
ProcessLaunchCallback processLaunchCallback,
ProcessExitCallback processExitCallback)
{
AZ_TestImpact_Eval(!m_executionState, ProcessException, "Couldn't execute schedule, schedule already in progress");
m_executionState = AZStd::make_unique<ExecutionState>(
m_maxConcurrentProcesses, processTimeout, scheduleTimeout, processLaunchCallback, processExitCallback);
const auto result = m_executionState->MonitorProcesses(processes);
m_executionState.reset();
return result;
}
} // namespace TestImpact
@@ -22,6 +22,7 @@
#include <AzCore/std/functional.h>
#include <AzCore/std/optional.h>
#include <AzCore/std/string/string.h>
#include <AzCore/std/smart_ptr/unique_ptr.h>
namespace TestImpact
{
@@ -49,6 +50,13 @@ namespace TestImpact
Abort //!< Abort scheduling immediately.
};
//! Result of the process scheduling sequence.
enum class ProcessSchedulerResult : bool
{
Graceful, //!< The scheduler completed its run without incident or was terminated gracefully in response to a client callback result.
Timeout //!< The scheduler aborted its run prematurely due to its runtime exceeding the scheduler timeout value.
};
//! Callback for process launch attempt.
//! @param processId The id of the process that attempted to launch.
//! @param launchResult The result of the process launch attempt.
@@ -79,38 +87,28 @@ namespace TestImpact
{
public:
//! Constructs the scheduler with the specified batch of processes.
//! @param processes The batch of processes to schedule.
//! @param processLaunchCallback The process launch callback function.
//! @param processExitCallback The process exit callback function.
//! @param maxConcurrentProcesses The maximum number of concurrent processes in-flight.
//! @param processTimeout The maximum duration a process may be in-flight for before being forcefully terminated.
//! @param scheduleTimeout The maximum duration the scheduler may run before forcefully terminating all in-flight processes.
//! processes and abandoning any queued processes.
ProcessScheduler(
const AZStd::vector<ProcessInfo>& processes,
const ProcessLaunchCallback& processLaunchCallback,
const ProcessExitCallback& processExitCallback,
size_t maxConcurrentProcesses,
AZStd::optional<AZStd::chrono::milliseconds> processTimeout,
AZStd::optional<AZStd::chrono::milliseconds> scheduleTimeout);
ProcessScheduler(size_t maxConcurrentProcesses);
~ProcessScheduler();
//! Executes the specified processes and calls the client callbacks (if any) as each process progresses in its life cycle.
//! @note Multiple subsequent calls to Execute are permitted.
//! @param processes The batch of processes to schedule.
//! @param processTimeout The maximum duration a process may be in-flight for before being forcefully terminated.
//! @param scheduleTimeout The maximum duration the scheduler may run before forcefully terminating all in-flight processes.
//! @param processLaunchCallback The process launch callback function.
//! @param processExitCallback The process exit callback function.
//! @returns The state that triggered the end of the schedule sequence.
ProcessSchedulerResult Execute(
const AZStd::vector<ProcessInfo>& processes,
AZStd::optional<AZStd::chrono::milliseconds> processTimeout,
AZStd::optional<AZStd::chrono::milliseconds> scheduleTimeout,
ProcessLaunchCallback processLaunchCallback,
ProcessExitCallback processExitCallback);
private:
struct ProcessInFlight;
void MonitorProcesses();
ProcessCallbackResult PopAndLaunch(ProcessInFlight& processInFlight);
void TerminateAllProcesses(ExitCondition exitStatus);
StdContent ConsumeProcessStdContent(ProcessInFlight& processInFlight);
void AccumulateProcessStdContent(ProcessInFlight& processInFlight);
const ProcessLaunchCallback m_processCreateCallback;
const ProcessExitCallback m_processExitCallback;
const AZStd::optional<AZStd::chrono::milliseconds> m_processTimeout;
const AZStd::optional<AZStd::chrono::milliseconds> m_scheduleTimeout;
const AZStd::chrono::high_resolution_clock::time_point m_startTime;
AZStd::vector<ProcessInFlight> m_processPool;
AZStd::queue<ProcessInfo> m_processQueue;
class ExecutionState;
AZStd::unique_ptr<ExecutionState> m_executionState;
size_t m_maxConcurrentProcesses = 0;
};
} // namespace TestImpact