Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
173 changes: 155 additions & 18 deletions src/sim/scheduler_fcfs_custom.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,28 +19,41 @@ CustomFCFSScheduler::CustomFCFSScheduler(
size_t num_max_candidates, job_cost_function_t cost_function,
backfill_selector_t selector, size_t initial_capacity,
CircularOverflowPolicy overflow_policy)
: CustomFCFSScheduler(total_nodes, initial_job_count, bf_policy,
num_max_candidates, initial_capacity,
overflow_policy) {
m_job_cost_function = std::move(cost_function);
m_backfill_selector = std::move(selector);
if (!m_job_cost_function) {
throw std::invalid_argument(
"CustomFCFSScheduler requires a job cost function");
}
if (!m_backfill_selector) {
throw std::invalid_argument(
"CustomFCFSScheduler requires a backfill selector");
}
}

CustomFCFSScheduler::CustomFCFSScheduler(num_nodes_t total_nodes,
size_t initial_job_count,
BackfillPolicy bf_policy,
size_t num_max_candidates,
size_t initial_capacity,
CircularOverflowPolicy overflow_policy)
: SchedulerBase(total_nodes, bf_policy),
m_wait_queue(initial_capacity != 0 ? initial_capacity
: initial_job_count),
m_overflow_policy(overflow_policy), m_eligible_end_idx(0),
m_current_tracked_time(0.0), m_removed_count(0),
m_num_max_candidates(num_max_candidates),
m_job_cost_function(std::move(cost_function)),
m_backfill_selector(std::move(selector)), m_resource_area(0.0),
m_num_max_candidates(num_max_candidates), m_newly_eligible_begin_idx(0),
m_arrival_update_pending(false), m_reevaluate_all_candidates(false),
m_candidate_preparation_pending(true), m_resource_area(0.0),
m_resource_area_time(0.0), m_resource_area_start(0.0),
m_accounted_available_nodes(total_nodes) {
if (m_num_max_candidates == 0) {
throw std::invalid_argument(
"CustomFCFSScheduler requires num_max_candidates > 0");
}
if (!m_job_cost_function) {
throw std::invalid_argument(
"CustomFCFSScheduler requires a job cost function");
}
if (!m_backfill_selector) {
throw std::invalid_argument(
"CustomFCFSScheduler requires a backfill selector");
}
if (m_backfill_policy == BackfillPolicy::CONSERVATIVE) {
throw std::invalid_argument(
"CustomFCFSScheduler supports EASY or NONE backfilling");
Expand Down Expand Up @@ -168,9 +181,73 @@ CustomFCFSScheduler::prediction_horizon(const running_jobs_t &running_jobs,
(effective_utilization * static_cast<double>(m_total_nodes));
}

std::optional<job_no_t> CustomFCFSScheduler::select_backfill_candidate(
const backfill_candidates_t &candidates, num_nodes_t available_nodes,
const running_jobs_t &effective_running_jobs, sim_time_t current_time) {
(void)available_nodes;
(void)effective_running_jobs;
(void)current_time;
if (!m_backfill_selector) {
throw std::logic_error(
"CustomFCFSScheduler subclass did not implement candidate selection");
}
return m_backfill_selector(candidates);
}

void CustomFCFSScheduler::on_jobs_became_eligible(
size_t newly_eligible_begin, size_t eligible_end,
num_nodes_t available_nodes, const running_jobs_t &running_jobs,
sim_time_t current_time) {
(void)newly_eligible_begin;
(void)eligible_end;
(void)available_nodes;
(void)running_jobs;
(void)current_time;
}

void CustomFCFSScheduler::on_backfill_candidates_ready(
const backfill_candidates_t &candidates, num_nodes_t available_nodes,
const running_jobs_t &effective_running_jobs, sim_time_t current_time,
bool fcfs_jobs_started) {
(void)candidates;
(void)available_nodes;
(void)effective_running_jobs;
(void)current_time;
(void)fcfs_jobs_started;
}

void CustomFCFSScheduler::on_scheduling_cycle_complete(
num_nodes_t available_nodes, const running_jobs_t &running_jobs,
sim_time_t current_time) {
(void)available_nodes;
(void)running_jobs;
(void)current_time;
}

void CustomFCFSScheduler::finish_scheduling_cycle(
num_nodes_t available_nodes, const running_jobs_t &running_jobs,
sim_time_t current_time) {
on_scheduling_cycle_complete(available_nodes, running_jobs, current_time);
complete_candidate_scan();
}

void CustomFCFSScheduler::insert_job(job_no_t job_id, sim_time_t submit_time,
tdiff_t run_time_estimate,
num_nodes_t nodes_requested) {
if (!m_job_cost_function) {
throw std::logic_error(
"CustomFCFSScheduler subclass must use insert_job_with_cost");
}
insert_job_with_cost(job_id, submit_time, run_time_estimate, nodes_requested,
m_job_cost_function(job_id, submit_time,
run_time_estimate, nodes_requested));
}

void CustomFCFSScheduler::insert_job_with_cost(job_no_t job_id,
sim_time_t submit_time,
tdiff_t run_time_estimate,
num_nodes_t nodes_requested,
job_cost_t cost) {
if (m_wait_queue.full()) {
if (m_overflow_policy == CircularOverflowPolicy::ABORT) {
throw std::runtime_error("CustomFCFSScheduler: wait queue capacity (" +
Expand All @@ -180,23 +257,31 @@ void CustomFCFSScheduler::insert_job(job_no_t job_id, sim_time_t submit_time,
m_wait_queue.set_capacity(std::max<size_t>(m_wait_queue.capacity() * 2, 1));
}

const job_cost_t cost = m_job_cost_function(
job_id, submit_time, run_time_estimate, nodes_requested);
m_wait_queue.push_back(
JobEntry(job_id, submit_time, run_time_estimate, nodes_requested, cost));
if (submit_time <= m_current_tracked_time) {
m_newly_eligible_begin_idx =
std::min(m_newly_eligible_begin_idx, m_eligible_end_idx);
m_eligible_end_idx = m_wait_queue.size();
m_arrival_update_pending = true;
m_candidate_preparation_pending = true;
}
}

void CustomFCFSScheduler::sync_to(sim_time_t current_time) {
if (current_time <= m_current_tracked_time) {
return;
}
m_newly_eligible_begin_idx = m_eligible_end_idx;
const size_t previous_eligible_end = m_eligible_end_idx;
while (m_eligible_end_idx < m_wait_queue.size() &&
m_wait_queue[m_eligible_end_idx].submit_time <= current_time) {
++m_eligible_end_idx;
}
if (m_eligible_end_idx != previous_eligible_end) {
m_arrival_update_pending = true;
m_candidate_preparation_pending = true;
}
m_current_tracked_time = current_time;
}

Expand All @@ -215,20 +300,34 @@ void CustomFCFSScheduler::compact_if_needed() {
}

size_t eligible_removed = 0;
size_t newly_eligible_removed = 0;
for (size_t i = 0; i < m_eligible_end_idx; ++i) {
eligible_removed += m_wait_queue[i].removed ? 1 : 0;
if (m_wait_queue[i].removed) {
++eligible_removed;
if (i < m_newly_eligible_begin_idx) {
++newly_eligible_removed;
}
}
}
const auto new_end =
std::remove_if(m_wait_queue.begin(), m_wait_queue.end(),
[](const JobEntry &entry) { return entry.removed; });
m_wait_queue.erase(new_end, m_wait_queue.end());
m_eligible_end_idx -= eligible_removed;
m_newly_eligible_begin_idx -= newly_eligible_removed;
m_removed_count = 0;
}

backfill_candidates_t CustomFCFSScheduler::find_backfill_candidates(
num_nodes_t available_nodes, sim_time_t current_time,
sim_time_t reservation_time) const {
return find_backfill_candidates_from(1, available_nodes, current_time,
reservation_time);
}

backfill_candidates_t CustomFCFSScheduler::find_backfill_candidates_from(
size_t begin_idx, num_nodes_t available_nodes, sim_time_t current_time,
sim_time_t reservation_time) const {
backfill_candidates_t candidates;
candidates.reserve(std::min(m_num_max_candidates, m_eligible_end_idx > 0
? m_eligible_end_idx - 1
Expand All @@ -238,7 +337,7 @@ backfill_candidates_t CustomFCFSScheduler::find_backfill_candidates(
return candidates;
}

for (size_t i = 1; i < m_eligible_end_idx; ++i) {
for (size_t i = std::max<size_t>(1, begin_idx); i < m_eligible_end_idx; ++i) {
const auto &job = m_wait_queue[i];
if (job.removed || job.nodes_requested > available_nodes) {
continue;
Expand All @@ -261,7 +360,14 @@ CustomFCFSScheduler::schedule(num_nodes_t free_nodes,
sync_to(current_time);
compact_if_needed();

if (m_arrival_update_pending) {
on_jobs_became_eligible(m_newly_eligible_begin_idx, m_eligible_end_idx,
free_nodes, running_jobs, current_time);
m_arrival_update_pending = false;
}

if (m_eligible_end_idx == 0) {
finish_scheduling_cycle(free_nodes, running_jobs, current_time);
return {};
}

Expand All @@ -283,24 +389,55 @@ CustomFCFSScheduler::schedule(num_nodes_t free_nodes,
}
m_wait_queue.pop_front();
--m_eligible_end_idx;
if (m_newly_eligible_begin_idx > 0) {
--m_newly_eligible_begin_idx;
}
}

if (active_job_count() == 0 || m_backfill_policy == BackfillPolicy::NONE) {
if (jobs_to_run.empty()) {
finish_scheduling_cycle(available_nodes, effective_running_jobs,
current_time);
} else {
complete_candidate_scan();
}
return jobs_to_run;
}

m_fcfs_reservation_time = calculate_fcfs_reservation(
m_wait_queue.front().nodes_requested, available_nodes,
effective_running_jobs, current_time);

const auto candidates = find_backfill_candidates(
available_nodes, current_time, m_fcfs_reservation_time);
const size_t candidate_begin =
m_reevaluate_all_candidates ? 1 : m_newly_eligible_begin_idx;
const auto candidates = find_backfill_candidates_from(
candidate_begin, available_nodes, current_time, m_fcfs_reservation_time);
if (candidates.empty()) {
if (jobs_to_run.empty()) {
finish_scheduling_cycle(available_nodes, effective_running_jobs,
current_time);
} else {
complete_candidate_scan();
}
return jobs_to_run;
}

const std::optional<job_no_t> selected = m_backfill_selector(candidates);
if (m_candidate_preparation_pending) {
on_backfill_candidates_ready(candidates, available_nodes,
effective_running_jobs, current_time,
!jobs_to_run.empty());
m_candidate_preparation_pending = false;
}

const std::optional<job_no_t> selected = select_backfill_candidate(
candidates, available_nodes, effective_running_jobs, current_time);
if (!selected) {
if (jobs_to_run.empty()) {
finish_scheduling_cycle(available_nodes, effective_running_jobs,
current_time);
} else {
complete_candidate_scan();
}
return jobs_to_run;
}

Expand Down
83 changes: 83 additions & 0 deletions src/sim/scheduler_fcfs_custom.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ class CustomFCFSScheduler : public SchedulerBase {
private:
template <typename TraceType> friend class BasicSimulation;

protected:
/** Queue entry exposed read-only to scheduler subclasses. */
struct JobEntry {
job_no_t job_id;
sim_time_t submit_time;
Expand All @@ -56,12 +58,17 @@ class CustomFCFSScheduler : public SchedulerBase {
nodes_requested(nodes), m_cost(cost), removed(false) {}
};

private:
boost::circular_buffer<JobEntry> m_wait_queue;
CircularOverflowPolicy m_overflow_policy;
size_t m_eligible_end_idx;
sim_time_t m_current_tracked_time;
size_t m_removed_count;
size_t m_num_max_candidates;
size_t m_newly_eligible_begin_idx;
bool m_arrival_update_pending;
bool m_reevaluate_all_candidates;
bool m_candidate_preparation_pending;
job_cost_function_t m_job_cost_function;
backfill_selector_t m_backfill_selector;

Expand Down Expand Up @@ -101,6 +108,60 @@ class CustomFCFSScheduler : public SchedulerBase {
tdiff_t prediction_horizon(const running_jobs_t &running_jobs,
sim_time_t current_time, double utilization) const;

/**
* Construct the FCFS scheduling core for a subclass-owned selection policy.
* The subclass must override select_backfill_candidate() and may insert
* precomputed queue costs with insert_job_with_cost().
*/
CustomFCFSScheduler(
num_nodes_t total_nodes, size_t initial_job_count,
BackfillPolicy bf_policy, size_t num_max_candidates,
size_t initial_capacity,
CircularOverflowPolicy overflow_policy = CircularOverflowPolicy::GROW);

/** Insert a job with a cost computed by a subclass-owned policy. */
void insert_job_with_cost(job_no_t job_id, sim_time_t submit_time,
tdiff_t run_time_estimate,
num_nodes_t nodes_requested, job_cost_t cost);

/** Return a read-only view of all stored queue entries. */
const boost::circular_buffer<JobEntry> &queued_jobs() const {
return m_wait_queue;
}

/** Return the exclusive end of the currently eligible queue range. */
size_t eligible_job_end() const { return m_eligible_end_idx; }

/** Allow a subclass to choose from the bounded feasible candidate set. */
virtual std::optional<job_no_t> select_backfill_candidate(
const backfill_candidates_t &candidates, num_nodes_t available_nodes,
const running_jobs_t &effective_running_jobs, sim_time_t current_time);

/**
* Called once after a complete same-timestamp arrival batch becomes
* eligible and before any FCFS dispatch from that batch.
*/
virtual void on_jobs_became_eligible(size_t newly_eligible_begin,
size_t eligible_end,
num_nodes_t available_nodes,
const running_jobs_t &running_jobs,
sim_time_t current_time);

/**
* Called once after FCFS dispatch and before the first backfill selection
* for an event cycle.
*/
virtual void
on_backfill_candidates_ready(const backfill_candidates_t &candidates,
num_nodes_t available_nodes,
const running_jobs_t &effective_running_jobs,
sim_time_t current_time, bool fcfs_jobs_started);

/** Called after no additional job can be dispatched at the current time. */
virtual void on_scheduling_cycle_complete(num_nodes_t available_nodes,
const running_jobs_t &running_jobs,
sim_time_t current_time);

public:
/**
* @param[in] total_nodes Cluster capacity available for allocations.
Expand Down Expand Up @@ -154,6 +215,28 @@ class CustomFCFSScheduler : public SchedulerBase {
size_t wait_queue_size() const override { return m_wait_queue.size(); }

private:
/** Request a full candidate scan after completions or capacity changes. */
void notify_resource_change() {
m_reevaluate_all_candidates = true;
m_candidate_preparation_pending = true;
}

/** Complete one event-time scheduling cycle and reset scan tracking. */
void finish_scheduling_cycle(num_nodes_t available_nodes,
const running_jobs_t &running_jobs,
sim_time_t current_time);

void complete_candidate_scan() {
m_newly_eligible_begin_idx = m_eligible_end_idx;
m_reevaluate_all_candidates = false;
m_candidate_preparation_pending = false;
}

backfill_candidates_t
find_backfill_candidates_from(size_t begin_idx, num_nodes_t available_nodes,
sim_time_t current_time,
sim_time_t reservation_time) const;

void compact_if_needed();
};

Expand Down
Loading
Loading