diff --git a/CHANGELOG.md b/CHANGELOG.md index c17ad77d7..fab54456a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,11 +5,13 @@ ### New features * The server scheduler now contains a safety limit for computation, configurable via `--scheduler-time-limit` (default: 5s) +* Better scheduling policy (prefill) for heterogenous clusters ### Fixes +* Fixed some occasional greedy backfilling in server scheduler + improvemnts in the reservation algorithm * Fixed server crash in a specific situation when an unschedulable high-priority task occurs - +* Fixed server crash caused by invalid handling of prefill ## v0.26.2 diff --git a/crates/tako/src/internal/scheduler/mapping.rs b/crates/tako/src/internal/scheduler/mapping.rs index 306110fc4..b5be85838 100644 --- a/crates/tako/src/internal/scheduler/mapping.rs +++ b/crates/tako/src/internal/scheduler/mapping.rs @@ -161,10 +161,15 @@ fn process_proactive_filling(core: &mut Core, mapping: &mut WorkerTaskMapping) { task_map, worker_map, task_queues, - request_map: _, + request_map, scheduler_state, .. } = core.split_mut(); + let max_prefill = scheduler_state.config.proactive_filling_max as u64; + if max_prefill == 0 { + // Prefill explicitly disabled. + return; + } let top_priority = task_queues.top_priority(); for queue in task_queues.iter_mut() { if queue.top_priority() != Some(top_priority) { @@ -176,6 +181,13 @@ fn process_proactive_filling(core: &mut Core, mapping: &mut WorkerTaskMapping) { if size == 0 { continue; } + let rqv = request_map.get(queue.resource_rq_id); + let max_capacity = worker_map + .get_workers() + .map(|w| w.resources.task_max_count(rqv)) + .max() + .unwrap_or(1) + .max(1) as u64; let workers: Vec<_> = worker_map .values_mut() .filter(|worker| { @@ -207,28 +219,45 @@ fn process_proactive_filling(core: &mut Core, mapping: &mut WorkerTaskMapping) { if workers.is_empty() { continue; } - let prefill_size = - (size / workers.len() as u32).min(scheduler_state.config.proactive_filling_max); - if prefill_size == 0 { - continue; - } - for worker in workers { - let tasks = queue.take_tasks_for_prefill(prefill_size); - for task_id in &tasks { + let capacities: Vec = workers + .iter() + .map(|w| w.resources.task_max_count(rqv).max(1) as u64) + .collect(); + let total_capacity: u64 = capacities.iter().sum(); + let mut return_back = Vec::new(); + for (worker, capacity) in workers.into_iter().zip(capacities) { + // The shares sum to at most `size`, so the queue entry we are drawing from is + // never exhausted before the last worker. + let share = size as u64 * capacity / total_capacity; + let depth = (max_prefill * capacity / max_capacity).max(1); + let prefill_size = share.min(depth); + if prefill_size == 0 { + continue; + } + let prefills = &mut mapping.workers.entry(worker.id).or_default().prefills; + for _ in 0..prefill_size { + let task_id = queue.take_one().unwrap(); log::debug!("Prefiling task={task_id} to worker={}", worker.id); - let task = task_map.get_task_mut(*task_id); - assert!(task.is_waiting()); - task.state = TaskRuntimeState::Prefilled { - worker_id: worker.id, - }; - worker.insert_prefill_task(*task_id); + let task = task_map.get_task_mut(task_id); + if task.is_waiting() { + task.state = TaskRuntimeState::Prefilled { + worker_id: worker.id, + }; + worker.insert_prefill_task(task_id); + queue.insert_prefill(task_id, top_priority, prefill_size as usize); + prefills.push(task_id); + } else { + // This can happen when task is in retracting, and it should be queite rare + log::debug!( + "Task is not in waiting state ({:?}) back to the queue.", + task.state + ); + return_back.push(task_id); + } } - mapping - .workers - .entry(worker.id) - .or_default() - .prefills - .extend(tasks); + } + for task_id in return_back { + queue.return_back(task_id, top_priority); } } } diff --git a/crates/tako/src/internal/scheduler/solver.rs b/crates/tako/src/internal/scheduler/solver.rs index 719df80e3..91abb24f6 100644 --- a/crates/tako/src/internal/scheduler/solver.rs +++ b/crates/tako/src/internal/scheduler/solver.rs @@ -1,7 +1,8 @@ -use crate::internal::common::resources::{ResourceId, ResourceRequest}; +use crate::internal::common::resources::{ResourceId, ResourceRequest, ResourceRequestVariants}; use crate::internal::scheduler::TaskBatch; use crate::internal::server::core::{Core, CoreSplit}; use crate::internal::server::worker::Worker; +use crate::internal::server::workerload::WorkerResources; use crate::internal::solver::{ConstraintType, LpSolution, LpSolver, Variable}; use crate::resources::{CPU_RESOURCE_ID, ResourceRqId}; use crate::{Map, ResourceVariantId, Set, WorkerId}; @@ -33,6 +34,11 @@ impl SchedulingSolution { } } +struct Candidates { + demand: u32, + candidates: Vec<(f64, WorkerId)>, +} + pub(crate) fn run_scheduling_solver( core: &Core, now: std::time::Instant, @@ -47,7 +53,7 @@ pub(crate) fn run_scheduling_solver( task_queues: _, request_map, worker_groups, - scheduler_state: scheduler_cache, + scheduler_state, .. } = core.split(); if request_map.is_empty() { @@ -91,6 +97,17 @@ pub(crate) fn run_scheduling_solver( let mut worker_res_constraint = vec![Vec::new(); n_resources]; let mut worker_cpu_constraint_no_reserves: Vec<(Variable, f64)> = Vec::with_capacity(task_batches.len()); + + let mut reservation_candidates: Map = Map::new(); + for batch in task_batches { + reservation_candidates.insert( + batch.resource_rq_id, + Candidates { + demand: batch.size, + candidates: Vec::new(), + }, + ); + } // Create worker-task placements for (w_idx, worker) in workers.iter().enumerate() { worker_cpu_constraint_no_reserves.clear(); @@ -150,13 +167,25 @@ pub(crate) fn run_scheduling_solver( } } + if has_variant + && !rqv.is_multi_node() + && batch.is_blocker + && let Some(a) = worker.sn_assignment() + { + let demand = &mut reservation_candidates + .get_mut(&batch.resource_rq_id) + .unwrap() + .demand; + *demand = demand.saturating_sub(a.free_resources.task_max_count(rqv)); + } + if !has_variant && !rqv.is_multi_node() && batch.is_blocker && worker.is_capable_to_run_rqv(rqv, now) && let Some(a) = worker.sn_assignment() { - let weight = w_idx as f64 / (n_workers * 100) as f64; + let weight = -((n_workers - w_idx) as f64 / (n_workers * 1024) as f64); solver.set_name(|| format!("R{}:{}", worker.id, batch.resource_rq_id)); let v = solver.add_bool_variable(weight); tasks_count_vars @@ -166,6 +195,13 @@ pub(crate) fn run_scheduling_solver( for (res_id, count) in a.free_resources.iter_pairs() { worker_res_constraint[res_id.as_usize()].push((v, count.as_f64())); } + // The blocker cannot run here now, but this worker could host it once it drains. + // How close it already is decides which workers are held for the blocker below. + reservation_candidates + .get_mut(&batch.resource_rq_id) + .unwrap() + .candidates + .push((fit_ratio(&a.free_resources, rqv), worker.id)); } } @@ -190,6 +226,31 @@ pub(crate) fn run_scheduling_solver( c.clear(); } } + let reserved: Map> = reservation_candidates + .into_iter() + .map( + |( + rq_id, + Candidates { + mut candidates, + demand, + }, + )| { + if demand > 0 { + candidates.sort_unstable_by(|(fit_a, w_a), (fit_b, w_b)| { + fit_b.total_cmp(fit_a).then(w_b.cmp(w_a)) + }); + } + let held = candidates + .into_iter() + .take(demand as usize) + .map(|(_fit, w_id)| w_id) + .collect(); + (rq_id, held) + }, + ) + .collect(); + let mut task_counts_per_group: Map<(ResourceRqId, &str), Variable> = Map::new(); let mut temp = Vec::new(); for batch in task_batches.iter() { @@ -253,6 +314,7 @@ pub(crate) fn run_scheduling_solver( }; let mut zero_cond = Vec::new(); + let mut zero_cond_reserved = Vec::new(); let mut blocked_by_unbounded: Set = Set::new(); for batch in task_batches.iter() { @@ -274,6 +336,7 @@ pub(crate) fn run_scheduling_solver( for cut in &batch.cuts { for (blocker_rq_id, blocking_size) in &cut.blockers { zero_cond.clear(); + zero_cond_reserved.clear(); let blocker_rqv = request_map.get(*blocker_rq_id); if batch_rqv.is_multi_node() { for (group_name, group) in worker_groups.iter() { @@ -292,22 +355,31 @@ pub(crate) fn run_scheduling_solver( if !w.is_capable_to_run_rqv(blocker_rqv, now) { continue; } - let gap = scheduler_cache.gap_cache.get_gap( + let gap = scheduler_state.gap_cache.get_gap( *blocker_rq_id, batch.resource_rq_id, &w.resources, sn_assignment.assigned_tasks.iter().map(|task_id| { let t = task_map.get_task(*task_id); - (t.resource_rq_id, t.rv_id().unwrap()) + ( + t.resource_rq_id, + t.assigned_placement(&scheduler_state.redirects).unwrap().1, + ) }), request_map, ); + // A worker held for this blocker gets the unconditional form: with no + // blocker-count term there is nothing a reservation variable can discharge. + let is_reserved = reserved + .get(blocker_rq_id) + .is_some_and(|ws| ws.contains(&w.id)); if gap > 0 { let vars = batch_rqv.variant_ids().filter_map(|v_id| { placements.get(&(w.id, batch.resource_rq_id, v_id)).copied() }); let cut_size = cut.size as f64; if let Some(s) = blocking_size + && !is_reserved && let Some(blocking_v) = get_bvar(&mut solver, *blocker_rq_id, *s) { solver.set_name(|| { @@ -324,7 +396,7 @@ pub(crate) fn run_scheduling_solver( blocking_v, batch_size, ); - } else if blocking_size.is_none() { + } else if blocking_size.is_none() || is_reserved { solver.set_name(|| { format!( "w{}: limit #rq{} to {} + {} (gap) where it can run with rq{blocker_rq_id}", @@ -343,6 +415,9 @@ pub(crate) fn run_scheduling_solver( placements.get(&(w.id, batch.resource_rq_id, v_id)) { zero_cond.push(*var); + if is_reserved { + zero_cond_reserved.push(*var); + } } } } @@ -390,6 +465,21 @@ pub(crate) fn run_scheduling_solver( }*/ } } + // Workers held for the blocker take only what the priority rule allows anyway + // (`cut.size`), with no blocker-count escape, so their capacity accumulates. + if !zero_cond_reserved.is_empty() { + solver.set_name(|| { + format!( + "limit #rq{} to {} on workers held for rq{blocker_rq_id}", + batch.resource_rq_id, cut.size + ) + }); + solver.add_constraint( + ConstraintType::Max, + cut.size as f64, + zero_cond_reserved.iter().map(|v| (*v, 1.0)), + ); + } if zero_cond.is_empty() { continue; } @@ -430,7 +520,7 @@ pub(crate) fn run_scheduling_solver( } let mut result = SchedulingSolution::default(); - let Some((solution, is_optimal)) = solver.solve_bounded(scheduler_cache.config.mip_time_limit) + let Some((solution, is_optimal)) = solver.solve_bounded(scheduler_state.config.mip_time_limit) else { return result; }; @@ -539,6 +629,26 @@ fn add_min_utilization( worker_res_constraint.pop(); } +fn fit_ratio(free: &WorkerResources, rqv: &ResourceRequestVariants) -> f64 { + rqv.requests() + .iter() + .map(|rq| { + rq.entries() + .iter() + .map(|e| { + let available = free.get(e.resource_id).total_fractions() as f64; + match e.request.amount_or_none_if_all() { + // `all` is only satisfied by an untouched resource, so anything less is 0. + None => 0.0, + Some(required) if required.is_zero() => 1.0, + Some(required) => (available / required.total_fractions() as f64).min(1.0), + } + }) + .fold(f64::INFINITY, f64::min) + }) + .fold(0.0, f64::max) +} + fn create_sn_var( solver: &mut LpSolver, rq: &ResourceRequest, diff --git a/crates/tako/src/internal/scheduler/taskqueue.rs b/crates/tako/src/internal/scheduler/taskqueue.rs index e63a14c15..b58df48e5 100644 --- a/crates/tako/src/internal/scheduler/taskqueue.rs +++ b/crates/tako/src/internal/scheduler/taskqueue.rs @@ -151,6 +151,11 @@ impl TaskQueue { } } + #[inline] + pub fn return_back(&mut self, task_id: TaskId, priority: Priority) { + self.add(task_id, priority); + } + fn add(&mut self, task_id: TaskId, priority: Priority) { match self.queue.entry(Reverse(priority)) { Entry::Vacant(e) => { @@ -216,6 +221,36 @@ impl TaskQueue { } } + /// Removes the task if it is in the queue and returns whether it was removed. + /// A retracting task may have been already taken out of the queue by the scheduler. + pub fn remove_if_queued(&mut self, task_id: TaskId, priority: Priority) -> bool { + if let Some((p, ts)) = &mut self.prefill + && priority == *p + && ts.remove(&task_id) + { + return true; + } + match self.queue.entry(Reverse(priority)) { + Entry::Vacant(_) => false, + Entry::Occupied(mut e) => match e.get_mut() { + OneOrMoreTaskIds::One(v) => { + if *v != task_id { + return false; + } + e.remove(); + true + } + OneOrMoreTaskIds::More(tasks) => { + let found = tasks.remove(&task_id); + if tasks.is_empty() { + e.remove(); + } + found + } + }, + } + } + pub fn shrink_to_fit(&mut self) { // Do nothing } @@ -301,20 +336,14 @@ impl TaskQueue { } } - pub fn take_tasks_for_prefill(&mut self, mut count: u32) -> Vec { - let entry = self.queue.first_entry().unwrap(); - let mut result = Vec::with_capacity(count as usize); - let priority = entry.key().0; - take_from_entry(entry, &mut count, &mut result); + pub fn insert_prefill(&mut self, task_id: TaskId, priority: Priority, max_prefill: usize) { if let Some(prefill) = &mut self.prefill { - assert_eq!(prefill.0, priority); - for task_id in &result { - prefill.1.insert(*task_id); - } + prefill.1.insert(task_id); } else { - self.prefill = Some((priority, result.iter().copied().collect())) + let mut v = Set::with_capacity(max_prefill); + v.insert(task_id); + self.prefill = Some((priority, v)) } - result } pub fn take_tasks(&mut self, mut count: u32) -> Vec { diff --git a/crates/tako/src/internal/server/reactor.rs b/crates/tako/src/internal/server/reactor.rs index 66f30eaa8..7e34c0843 100644 --- a/crates/tako/src/internal/server/reactor.rs +++ b/crates/tako/src/internal/server/reactor.rs @@ -320,13 +320,20 @@ fn task_running( // By removing redirections first, we unassign the task so we can later assign it back // In theory, we could optimize this special case by doing nothing, but it should be quite rare // So I prefer to keep the code simple. - try_remove_redirection( + if !try_remove_redirection( worker_map, scheduler_state, request_map, task_id, task.resource_rq_id, - ); + ) { + // We have tried to retract the task, but it has already started. + // Without a redirection, the task is still in the ready queue and it has + // to be removed, otherwise the queue keeps an id of an already removed task. + task_queues + .get_mut(task.resource_rq_id) + .remove_if_queued(task.id, task.priority()); + } let rqv = request_map.get(task.resource_rq_id); worker_map .get_worker_mut(worker_id) @@ -546,13 +553,17 @@ fn task_finished( } TaskRuntimeState::Retracting { worker_id: w_id } => { assert_eq!(*w_id, worker_id); - try_remove_redirection( + if !try_remove_redirection( worker_map, scheduler_state, request_map, task_id, task.resource_rq_id, - ); + ) { + task_queues + .get_mut(task.resource_rq_id) + .remove_if_queued(task.id, task.priority()); + } } TaskRuntimeState::Prefilled { .. } | TaskRuntimeState::Waiting { .. } @@ -589,17 +600,23 @@ fn task_finished( true } +/// Returns true if a redirection was found (and removed). +/// A task with a redirection was already taken out of the ready queue by the scheduler, +/// a retracting task without a redirection is still waiting in the ready queue. fn try_remove_redirection( worker_map: &mut WorkerMap, scheduler_state: &mut SchedulerState, request_map: &ResourceRqMap, task_id: TaskId, resource_rq_id: ResourceRqId, -) { +) -> bool { if let Some((worker_id, rv_id)) = scheduler_state.redirects.remove(&task_id) { let worker = worker_map.get_worker_mut(worker_id); let rq = request_map.get(resource_rq_id).get(rv_id); worker.remove_sn_task(task_id, rq); + true + } else { + false } } @@ -653,13 +670,17 @@ fn task_failed( } TaskRuntimeState::Retracting { worker_id: w } => { assert_eq!(worker_id, *w); - try_remove_redirection( + if !try_remove_redirection( worker_map, scheduler_state, request_map, task_id, task.resource_rq_id, - ); + ) { + task_queues + .get_mut(task.resource_rq_id) + .remove_if_queued(task.id, task.priority()); + } } _ => {} } diff --git a/crates/tako/src/internal/server/task.rs b/crates/tako/src/internal/server/task.rs index 35ccd9fc2..f180cdaef 100644 --- a/crates/tako/src/internal/server/task.rs +++ b/crates/tako/src/internal/server/task.rs @@ -247,6 +247,7 @@ impl Task { } } + #[cfg(test)] pub(crate) fn rv_id(&self) -> Option { match self.state { TaskRuntimeState::Running { rv_id, .. } | TaskRuntimeState::Assigned { rv_id, .. } => { @@ -256,6 +257,19 @@ impl Task { } } + #[inline] + pub(crate) fn assigned_placement( + &self, + redirects: &Map, + ) -> Option<(WorkerId, ResourceVariantId)> { + match self.state { + TaskRuntimeState::Assigned { worker_id, rv_id } + | TaskRuntimeState::Running { worker_id, rv_id } => Some((worker_id, rv_id)), + TaskRuntimeState::Retracting { .. } => redirects.get(&self.id).copied(), + _ => None, + } + } + pub(crate) fn increment_instance_id(&mut self) { self.instance_id = InstanceId::new(self.instance_id.as_num() + 1); } diff --git a/crates/tako/src/internal/server/worker.rs b/crates/tako/src/internal/server/worker.rs index 923714530..9790fc22d 100644 --- a/crates/tako/src/internal/server/worker.rs +++ b/crates/tako/src/internal/server/worker.rs @@ -243,15 +243,9 @@ impl Worker { let mut resources = self.resources.clone(); for task_id in a.assigned_tasks.iter() { let task = task_map.get_task(*task_id); - let (worker_id, rv_id) = match &task.state { - TaskRuntimeState::Assigned { worker_id, rv_id } - | TaskRuntimeState::Running { worker_id, rv_id } => (*worker_id, *rv_id), - TaskRuntimeState::Retracting { .. } => { - let (worker_id, rv_id) = transfers.get(task_id).unwrap(); - (*worker_id, *rv_id) - } - s => panic!("Invalid state {s:?}"), - }; + let (worker_id, rv_id) = task + .assigned_placement(transfers) + .unwrap_or_else(|| panic!("Invalid state {:?}", task.state)); assert_eq!(self.id, worker_id); let rq = request_map.get(task.resource_rq_id).get(rv_id); assert!(resources.is_capable_to_run_request(rq)); diff --git a/crates/tako/src/internal/tests/test_reactor.rs b/crates/tako/src/internal/tests/test_reactor.rs index ddb87cf26..29c90d524 100644 --- a/crates/tako/src/internal/tests/test_reactor.rs +++ b/crates/tako/src/internal/tests/test_reactor.rs @@ -16,7 +16,6 @@ use crate::internal::tests::utils::sorted_vec; use crate::internal::tests::utils::task::{TaskBuilder, task_running_msg}; use crate::internal::tests::utils::workflows::{submit_example_1, submit_example_3}; use crate::internal::worker::configuration::OverviewConfiguration; -use crate::internal::worker::task::RunningTask; use crate::resources::{ResourceAmount, ResourceDescriptorItem, ResourceIdMap}; use crate::tests::utils::env::{TestComm, TestEnv}; use crate::tests::utils::worker::WorkerBuilder; @@ -772,6 +771,14 @@ fn test_task_reject3() { assert!(rt.task(t2).is_waiting()); } +fn get_prefilled(rt: &mut TestEnv, tasks: &[TaskId]) -> Option { + tasks.iter().find(|t| rt.task(**t).is_prefilled()).copied() +} + +fn get_assigned(rt: &mut TestEnv, tasks: &[TaskId]) -> Option { + tasks.iter().find(|t| rt.task(**t).is_assigned()).copied() +} + fn setup_prefill(rt: &mut TestEnv) -> (WorkerId, TaskId, TaskId) { rt.set_scheduler_config(SchedulerConfig { proactive_filling_reserve: 1, @@ -781,16 +788,8 @@ fn setup_prefill(rt: &mut TestEnv) -> (WorkerId, TaskId, TaskId) { let tasks = rt.new_tasks(3, &TaskBuilder::new()); let w1 = rt.new_worker(&WorkerBuilder::new(1)); rt.schedule(); - let prefilled = tasks - .iter() - .find(|t| rt.task(**t).is_prefilled()) - .copied() - .unwrap(); - let assigned = tasks - .iter() - .find(|t| rt.task(**t).is_assigned()) - .copied() - .unwrap(); + let prefilled = get_prefilled(rt, &tasks).unwrap(); + let assigned = get_assigned(rt, &tasks).unwrap(); (w1, assigned, prefilled) } @@ -825,6 +824,30 @@ fn test_prefill_submit_high_priority() { } } +#[test] +fn test_prefill_retracted_and_prefill_again() { + let mut rt = TestEnv::new(); + rt.set_scheduler_config(SchedulerConfig { + proactive_filling_reserve: 0, + proactive_filling_max: 1, + ..Default::default() + }); + let tasks1 = rt.new_tasks(3, &TaskBuilder::new()); + let w = rt.new_worker(&WorkerBuilder::new(1)); + rt.schedule(); + let p1 = get_prefilled(&mut rt, &tasks1).unwrap(); + let a1 = get_assigned(&mut rt, &tasks1).unwrap(); + let new_task = rt.new_task(&TaskBuilder::new().user_priority(1)); + assert!(matches!( + rt.task(p1).state, + TaskRuntimeState::Retracting { worker_id } if w == worker_id + )); + rt.finish_task(a1, w); + rt.schedule(); + assert!(rt.task(new_task).is_assigned()); + assert!(rt.task(p1).is_retracting()); +} + #[test] fn test_prefill_submit_same_priority() { for cpus in [1, 2] { @@ -878,29 +901,44 @@ fn test_prefill_started_on_same_worker() { assert!(rt.task(t1).is_assigned()); let tasks = rt.new_tasks(2, &TaskBuilder::new()); rt.schedule(); - let prefilled: TaskId = tasks - .iter() - .find(|t| rt.task(**t).is_prefilled()) - .copied() - .unwrap(); - let assigned: TaskId = tasks - .iter() - .find(|t| rt.task(**t).is_assigned()) - .copied() - .unwrap(); + let prefilled: TaskId = get_prefilled(&mut rt, &tasks).unwrap(); let up1 = WorkerTaskUpdate::Finished { task_id: t1 }; let mut comm = TestComm::new(); on_task_update(rt.core(), &mut comm, w1, smallvec![up1]); - rt.schedule(); - assert!(rt.task(prefilled).is_retracting()); - let up2 = WorkerTaskUpdate::Running(task_running_msg(prefilled)); on_task_update(rt.core(), &mut comm, w1, smallvec![up2]); rt.sanity_check(); } +#[test] +fn test_prefill_retracted_but_already_started() { + let mut rt = TestEnv::new(); + let (w1, t1, t2) = setup_prefill(&mut rt); + let mut comm = TestComm::new(); + + let t3 = TaskId::new(100.into(), 501.into()); + let task3 = TaskBuilder::new().user_priority(10).build(t3, rt.core()); + on_new_tasks(rt.core(), &mut comm, vec![task3]); + comm.check_need_scheduling(); + match &comm.take_worker_msgs(w1, 1)[0] { + ToWorkerMessage::RetractTasks(ts) => assert_eq!(ts.ids, vec![t2]), + _ => panic!("Invalid worker msg"), + } + comm.emptiness_check(); + assert!(rt.task(t2).is_retracting()); + + rt.finish_task(t1, w1); + let up = WorkerTaskUpdate::Running(task_running_msg(t2)); + on_task_update(rt.core(), &mut comm, w1, smallvec![up]); + comm.client.take_task_running(1); + comm.check_need_scheduling(); + comm.emptiness_check(); + rt.finish_task(t2, w1); + rt.schedule(); +} + #[test] fn test_prefill_started() { let mut rt = TestEnv::new(); diff --git a/crates/tako/src/internal/tests/test_scheduler_sn.rs b/crates/tako/src/internal/tests/test_scheduler_sn.rs index 5e8b63ba0..f90c229bb 100644 --- a/crates/tako/src/internal/tests/test_scheduler_sn.rs +++ b/crates/tako/src/internal/tests/test_scheduler_sn.rs @@ -1,3 +1,4 @@ +use crate::internal::common::resources::ResourceId; use crate::internal::messages::worker::ToWorkerMessage; use crate::internal::scheduler::{PriorityCut, SchedulerConfig, create_task_batches}; use crate::internal::server::reactor::on_retract_response; @@ -7,7 +8,7 @@ use crate::resources::ResourceRqId; use crate::tests::utils::env::{TestComm, TestEnv}; use crate::tests::utils::task::TaskBuilder; use crate::tests::utils::worker::WorkerBuilder; -use crate::{ResourceVariantId, WorkerId}; +use crate::{ResourceVariantId, TaskId, WorkerId}; use std::time::Duration; #[test] @@ -1305,6 +1306,160 @@ fn test_prefill_steal() { rt.sanity_check(); } +#[test] +fn test_gap_over_redirected_retracting_task() { + let mut rt = TestEnv::new(); + rt.set_scheduler_config(SchedulerConfig { + proactive_filling_reserve: 3, + proactive_filling_max: 6, + ..Default::default() + }); + + // -- Precondition 1: a redirected retracting task in w2's `assigned_tasks`. + // Same opening as `test_prefill_steal`: w1 prefills, then the larger w2 steals part of that + // prefill, so those tasks become `Retracting` (on w1) while being assigned to w2. + let w1 = rt.new_worker(&WorkerBuilder::new(1)); + let low = rt.new_tasks(9, &TaskBuilder::new()); + let low_rq_id = rt.task(low[0]).resource_rq_id; + rt.schedule(); + assert_eq!(prefill_count(&mut rt, w1), 5, "w1 did not prefill"); + let w2 = rt.new_worker(&WorkerBuilder::new(5)); + rt.schedule(); + + // Asserted rather than assumed: if prefill ever stops producing this state, the test must + // fail loudly instead of passing while reproducing nothing. + assert!( + !rt.core().split().scheduler_state.redirects.is_empty(), + "no redirect was created, so the state under test does not exist" + ); + let retracting_on_w2 = rt + .worker(w2) + .sn_assignment() + .unwrap() + .assigned_tasks + .iter() + .filter(|task_id| rt.task(**task_id).rv_id().is_none()) + .count(); + assert!( + retracting_on_w2 > 0, + "w2 holds no assigned task without an rv_id; the unwrap cannot be reached" + ); + + // -- Precondition 2: spare capacity for the low-priority request. + // The solver skips a batch that has no placement variables, and after the steal both w1 and + // w2 are full -- so without this the low-priority batch, and with it the entire + // priority-condition path, is never visited. A worker added now cannot undo the redirect + // already recorded above. 2 cpus is deliberately too narrow for the blocker below, so this + // worker contributes capacity without becoming a candidate for it. + rt.new_worker(&WorkerBuilder::new(2)); + + // -- Precondition 3: a blocker, so the solver builds a priority condition at all. + // A cut is only emitted when a *second* queue is non-empty at a higher priority, so this + // needs a distinct resource request. 5 cpus fits w2's total width -- w2 is therefore + // `is_capable_to_run_rqv` and is not skipped by the impossible-filter -- but w2 has no room + // for it right now, which is what makes it block. + rt.new_tasks(2, &TaskBuilder::new().cpus(5).user_priority(10)); + + let batches = create_task_batches(rt.core(), std::time::Instant::now(), None); + let low_batch = batches + .iter() + .find(|b| b.resource_rq_id == low_rq_id) + .expect("the low-priority request must still have a batch"); + assert!( + low_batch.cuts.iter().any(|cut| !cut.blockers.is_empty()), + "no blocker cut was produced, so the gap path is never entered: {:?}", + low_batch.cuts + ); + + rt.schedule(); + rt.sanity_check(); +} + +/// Prefill has to be weighted by worker capacity, otherwise a small worker is handed as many +/// tasks as a large one and takes proportionally longer to drain them. Prefilled tasks are +/// removed from the global queue and are not reclaimed while regular tasks of the same priority +/// remain, so the small worker ends up sitting on an older job's tail long after every large +/// worker has moved on to newer jobs. +#[test] +fn test_prefill_weighted_by_worker_capacity() { + let mut rt = TestEnv::new(); + rt.set_scheduler_config(SchedulerConfig { + proactive_filling_reserve: 0, + proactive_filling_max: 32, + ..Default::default() + }); + let w_big = rt.new_worker(&WorkerBuilder::new(16)); + let w_small = rt.new_worker(&WorkerBuilder::new(2)); + rt.new_tasks(300, &TaskBuilder::new()); + rt.schedule(); + + // `proactive_filling_max` applies to the largest worker; everyone else is scaled down by + // capacity. An unweighted split would give both workers 32. + let big = prefill_count(&mut rt, w_big); + let small = prefill_count(&mut rt, w_small); + assert_eq!(big, 32); + assert_eq!(small, 4); + + // The property that actually matters: both hold the same *duration* of backlog, i.e. the + // same number of task generations (two each here). + assert_eq!(big / 16, small / 2); + rt.sanity_check(); +} + +#[test] +fn test_priority_is_not_overridden_by_weight() { + let narrow = TaskBuilder::new().cpus(1).user_priority(10); + let wide = TaskBuilder::new().cpus(4).user_priority(0).weight(10.0); + + let mut c = TestCase::new(); + c.n_tasks(10, &narrow); + c.n_tasks(4, &wide); + // All capacity goes to the high-priority 1-cpu tasks despite the 10x weight on the others. + c.w(&WorkerBuilder::new(8)).expect_request(8, &narrow); + c.w(&WorkerBuilder::new(2)).expect_request(2, &narrow); + c.check(); +} + +#[test] +fn test_weight_prefers_request_at_equal_priority() { + let narrow = TaskBuilder::new().cpus(1); + let wide = TaskBuilder::new().cpus(4).weight(4.0); + + let mut c = TestCase::new(); + c.n_tasks(10, &narrow); + c.n_tasks(2, &wide); + c.w(&WorkerBuilder::new(8)).expect_request(2, &wide); + c.w(&WorkerBuilder::new(2)).expect_request(2, &narrow); + c.check(); +} + +/// The prefill depth must be measured against the largest worker in the *cluster*, not the +/// largest one eligible for prefill in this round. A worker is skipped while it still holds +/// prefill of the request, so the eligible set is routinely all-small -- and if the reference +/// capacity is taken from it, the depth springs back to `proactive_filling_max` for a tiny +/// worker, which is the whole bug. +#[test] +fn test_prefill_depth_when_large_worker_is_ineligible() { + let mut rt = TestEnv::new(); + rt.set_scheduler_config(SchedulerConfig { + proactive_filling_reserve: 0, + proactive_filling_max: 32, + ..Default::default() + }); + let w_big = rt.new_worker(&WorkerBuilder::new(16)); + rt.new_tasks(400, &TaskBuilder::new()); + rt.schedule(); + assert_eq!(prefill_count(&mut rt, w_big), 32); + + // w_big now holds prefill of this request, so it is excluded from further prefill and + // only the 2-cpu worker is eligible. + let w_small = rt.new_worker(&WorkerBuilder::new(2)); + rt.schedule(); + assert_eq!(prefill_count(&mut rt, w_big), 32); + assert_eq!(prefill_count(&mut rt, w_small), 4); + rt.sanity_check(); +} + #[test] pub fn test_schedule_running() { let mut rt = TestEnv::new(); @@ -1511,3 +1666,332 @@ fn test_schedule_bounded_is_optimal_true_when_solve_converges() { assert!(rt.schedule_solution().is_optimal); } + +#[test] +fn test_schedule_reservation_priority() { + let mut c = TestCase::new(); + let ht = c.t(&TaskBuilder::new().cpus(6).user_priority(10)); + c.ts(6, &TaskBuilder::new().cpus(1)); + c.w(&WorkerBuilder::new(6)).running_c(6); + c.w(&WorkerBuilder::new(6)).expect_tasks(&[ht]); + c.check(); +} + +#[test] +fn test_schedule_reservation_used_when_worker_frees_up() { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(4, &WorkerBuilder::new(6)); + let running: Vec<_> = ws + .iter() + .map(|w| rt.new_task_running(&TaskBuilder::new().cpus(4), *w)) + .collect(); + let blocker = rt.new_task(&TaskBuilder::new().cpus(6).user_priority(10)); + let small = rt.new_tasks(30, &TaskBuilder::new().cpus(1)); + + fn on_worker(rt: &TestEnv, task_id: TaskId, worker_id: WorkerId) -> bool { + matches!( + rt.task(task_id).state, + TaskRuntimeState::Assigned { worker_id: w, .. } if w == worker_id + ) + } + + rt.schedule(); + + // One worker is held back for the blocker, the three others take two tasks each. + let reserved_idx = ws + .iter() + .position(|w| !small.iter().any(|t| on_worker(&rt, *t, *w))) + .expect("no worker was reserved for the blocker"); + assert_eq!( + small.iter().filter(|t| rt.task(**t).is_assigned()).count(), + 6 + ); + assert!(rt.task(blocker).is_waiting()); + + // The reserved worker becomes completely free, which is exactly what the blocker waits for. + let reserved = ws[reserved_idx]; + rt.finish_task(running[reserved_idx], reserved); + rt.schedule(); + + assert!( + on_worker(&rt, blocker, reserved), + "reserved worker {reserved} was not used for the blocker, state: {:?}", + rt.task(blocker).state + ); + assert_eq!( + small + .iter() + .filter(|t| on_worker(&rt, **t, reserved)) + .count(), + 0 + ); +} + +#[test] +fn test_schedule_lone_blocker_holds_freed_capacity() { + /// Narrow tasks placed while `n_wide` blockers are queued; must be 0 for every `n_wide`. + fn narrow_placed(n_wide: usize) -> usize { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(4, &WorkerBuilder::new(8)); + for (idx, w) in ws.iter().enumerate() { + // The first worker has just freed 2 of its 8 cpus; the rest are saturated. + let running = if idx == 0 { 6 } else { 8 }; + rt.new_task_running(&TaskBuilder::new().cpus(running), *w); + } + rt.new_tasks(n_wide, &TaskBuilder::new().cpus(8).user_priority(10)); + let narrow = rt.new_tasks(30, &TaskBuilder::new().cpus(1)); + rt.schedule(); + narrow.iter().filter(|t| rt.task(**t).is_assigned()).count() + } + + for n_wide in 1..=4 { + assert_eq!( + narrow_placed(n_wide), + 0, + "low priority work took the freed cpus with {n_wide} wide task(s) queued" + ); + } +} + +#[test] +fn test_schedule_lone_blocker_accumulates_capacity_over_rounds() { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(4, &WorkerBuilder::new(8)); + let mut running: Vec> = ws + .iter() + .map(|w| { + (0..8) + .map(|_| rt.new_task_running(&TaskBuilder::new().cpus(1), *w)) + .collect() + }) + .collect(); + let blocker = rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + rt.new_tasks(60, &TaskBuilder::new().cpus(1)); + + // Eight rounds of "one task finishes everywhere" is exactly what a held worker needs to reach + // 8 free cpus; the extra round gives the scheduler the chance to place the blocker afterwards. + for _tick in 0..=8 { + rt.schedule(); + if rt.task(blocker).is_assigned() { + return; + } + for (idx, w) in ws.iter().enumerate() { + if let Some(task_id) = running[idx].pop() { + rt.finish_task(task_id, *w); + } + } + } + + let free: Vec = ws + .iter() + .map(|w| { + let a = rt.worker(*w).sn_assignment().unwrap(); + format!("w{w}: {:?}", a.free_resources.get(ResourceId::new(0))) + }) + .collect(); + panic!( + "blocker never started; no worker reassembled 8 free cpus ({})", + free.join(", ") + ); +} + +#[test] +fn test_schedule_reservation_holds_emptiest_worker() { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(4, &WorkerBuilder::new(8)); + // w0 has 2 cpus free, w1 has 4, the rest are saturated. w1 is the closest to fitting an + // 8-cpu task, even though w0 comes first in worker order. + for (idx, w) in ws.iter().enumerate() { + let running = match idx { + 0 => 6, + 1 => 4, + _ => 8, + }; + rt.new_task_running(&TaskBuilder::new().cpus(running), *w); + } + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + let narrow = rt.new_tasks(30, &TaskBuilder::new().cpus(1)); + rt.schedule(); + + let placed_on = |rt: &TestEnv, worker_id: WorkerId| { + narrow + .iter() + .filter(|t| { + matches!(rt.task(**t).state, + TaskRuntimeState::Assigned { worker_id: w, .. } if w == worker_id) + }) + .count() + }; + assert_eq!(placed_on(&rt, ws[1]), 0, "the emptiest worker must be held"); + assert_eq!( + placed_on(&rt, ws[0]), + 2, + "the other worker must still be used" + ); +} + +#[test] +fn test_schedule_reservation_leaves_other_workers_for_backfill() { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(3, &WorkerBuilder::new(8)); + for w in &ws { + rt.new_task_running(&TaskBuilder::new().cpus(6), *w); + } + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + let narrow = rt.new_tasks(30, &TaskBuilder::new().cpus(1)); + rt.schedule(); + + let held: Vec<_> = ws + .iter() + .filter(|w| { + !narrow.iter().any(|t| { + matches!(rt.task(*t).state, + TaskRuntimeState::Assigned { worker_id, .. } if worker_id == **w) + }) + }) + .collect(); + assert_eq!(held.len(), 1, "exactly one worker is held for the blocker"); + assert_eq!( + narrow.iter().filter(|t| rt.task(**t).is_assigned()).count(), + 4, + "the two other workers keep their 2 free cpus busy" + ); +} + +#[test] +fn test_schedule_blockers_hold_all_needed_workers() { + const N_WORKERS: usize = 4; + + /// Narrow tasks placed with `n_partial` workers holding 2 freed cpus and `n_blockers` queued. + fn narrow_placed(n_partial: usize, n_blockers: usize) -> usize { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(N_WORKERS, &WorkerBuilder::new(8)); + for (idx, w) in ws.iter().enumerate() { + let running = if idx < n_partial { 6 } else { 8 }; + rt.new_task_running(&TaskBuilder::new().cpus(running), *w); + } + rt.new_tasks(n_blockers, &TaskBuilder::new().cpus(8).user_priority(10)); + let narrow = rt.new_tasks(30, &TaskBuilder::new().cpus(1)); + rt.schedule(); + narrow.iter().filter(|t| rt.task(**t).is_assigned()).count() + } + + for n_partial in 1..=N_WORKERS { + for n_blockers in 1..=N_WORKERS { + // The batch limit is one per capable worker, because no worker can host an 8-cpu task + // right now. At the limit `blocking_size` becomes `None`, the cut is unconditional on + // every capable worker and backfill stops there for pre-existing reasons that have + // nothing to do with holding. + let expected = if n_blockers >= N_WORKERS { + 0 + } else { + 2 * n_partial.saturating_sub(n_blockers) + }; + assert_eq!( + narrow_placed(n_partial, n_blockers), + expected, + "{n_partial} worker(s) with freed cpus, {n_blockers} blocker(s) queued" + ); + } + } +} + +#[test] +fn test_schedule_blockers_accumulate_in_parallel() { + /// Blockers started within one worker's drain; must be all of them. + fn started_within_one_drain(n_blockers: usize) -> (usize, Vec) { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(4, &WorkerBuilder::new(8)); + let mut running: Vec> = ws + .iter() + .map(|w| { + (0..8) + .map(|_| rt.new_task_running(&TaskBuilder::new().cpus(1), *w)) + .collect() + }) + .collect(); + let blockers = rt.new_tasks(n_blockers, &TaskBuilder::new().cpus(8).user_priority(10)); + rt.new_tasks(60, &TaskBuilder::new().cpus(1)); + + for _tick in 0..=8 { + rt.schedule(); + if blockers.iter().all(|t| rt.task(*t).is_assigned()) { + break; + } + for (idx, w) in ws.iter().enumerate() { + if let Some(task_id) = running[idx].pop() { + rt.finish_task(task_id, *w); + } + } + } + + let placed = blockers + .iter() + .filter(|t| rt.task(**t).is_assigned()) + .count(); + let free = ws + .iter() + .map(|w| { + let a = rt.worker(*w).sn_assignment().unwrap(); + format!("w{w}: {:?}", a.free_resources.get(ResourceId::new(0))) + }) + .collect(); + (placed, free) + } + + for n_blockers in 1..=4 { + let (placed, free) = started_within_one_drain(n_blockers); + assert_eq!( + placed, + n_blockers, + "only {placed} of {n_blockers} blockers started; workers drained one at a time ({})", + free.join(", ") + ); + } +} + +#[test] +fn test_schedule_wide_workers_held_one_per_blocker() { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(3, &WorkerBuilder::new(16)); + let mut running: Vec> = ws + .iter() + .map(|w| { + (0..16) + .map(|_| rt.new_task_running(&TaskBuilder::new().cpus(1), *w)) + .collect() + }) + .collect(); + let blockers = rt.new_tasks(2, &TaskBuilder::new().cpus(8).user_priority(10)); + rt.new_tasks(60, &TaskBuilder::new().cpus(1)); + + // Eight rounds is what a held worker needs to reach the 8 free cpus a blocker wants; both + // blockers must be served by that point, not one after the other. + for _tick in 0..=8 { + rt.schedule(); + if blockers.iter().all(|t| rt.task(*t).is_assigned()) { + return; + } + for (idx, w) in ws.iter().enumerate() { + if let Some(task_id) = running[idx].pop() { + rt.finish_task(task_id, *w); + } + } + } + + let placed = blockers + .iter() + .filter(|t| rt.task(**t).is_assigned()) + .count(); + let free: Vec = ws + .iter() + .map(|w| { + let a = rt.worker(*w).sn_assignment().unwrap(); + format!("w{w}: {:?}", a.free_resources.get(ResourceId::new(0))) + }) + .collect(); + panic!( + "only {placed} of 2 blockers started; one wide worker was credited for both ({})", + free.join(", ") + ); +} diff --git a/crates/tako/src/internal/tests/utils/scheduler.rs b/crates/tako/src/internal/tests/utils/scheduler.rs index 69ce4a7bc..167cad306 100644 --- a/crates/tako/src/internal/tests/utils/scheduler.rs +++ b/crates/tako/src/internal/tests/utils/scheduler.rs @@ -81,6 +81,14 @@ impl TestCase { self.rt.get_mut().new_tasks_cpus(cpus) } + /// `count` tasks from an explicit builder, for properties `pc_tasks` cannot express + /// (resource weight, variants, ...). + pub fn n_tasks(&mut self, count: usize, builder: &TaskBuilder) -> Vec { + (0..count) + .map(|_| self.rt.get_mut().new_task(builder)) + .collect() + } + // priority + cpu tasks pub fn pc_tasks(&mut self, priority_cpus: &[(i32, u32)]) -> Vec { priority_cpus diff --git a/tests/test_resources.py b/tests/test_resources.py index 01291aa05..a24ce2ddc 100644 --- a/tests/test_resources.py +++ b/tests/test_resources.py @@ -636,3 +636,64 @@ def test_scheduler_unschedulable_sn_blocker(hq_env: HqEnv): hq_env.check_running_processes() table = hq_env.command(["job", "info", "3"], as_table=True) assert table.get_row_value("State").endswith("WAITING (5)") + + +def test_scheduler_priority_churn(hq_env: HqEnv): + hq_env.start_server() + hq_env.start_workers(1, cpus=4) + hq_env.command(["submit", "--array=0-999", "--stdout=none", "--stderr=none", "--", "sleep", "0.05"]) + wait_for_job_state(hq_env, 1, "RUNNING") + + for priority in range(1, 6): + hq_env.command( + [ + "submit", + f"--priority={priority}", + "--stdout=none", + "--stderr=none", + "--", + "sleep", + "0.05", + ] + ) + time.sleep(0.5) + hq_env.check_running_processes() + + +def test_scheduler_reservation(hq_env: HqEnv, tmp_path): + hq_env.start_server() + hq_env.start_workers(4, cpus=6) + hq_env.command(["submit", "--array=1-4", "--cpus=4", "--", "sleep", "100"]) + wait_for_job_state(hq_env, 1, "RUNNING") + time.sleep(0.5) + content = [ + """ +[[task]] +id = 0 +priority = 10 +command = ["sleep", "100"] + +[[task.request]] +resources = { "cpus" = 6 } + """ + ] + for i in range(1, 31): + content.append(f""" +[[task]] +id = {i} +command = ["sleep", "100"] +[[task.request]] +resources = {{ "cpus" = 1 }} +""") + tmp_path.joinpath("job.toml").write_text("\n".join(content)) + hq_env.command(["job", "submit-file", "job.toml"]) + wait_for_job_state(hq_env, 1, "RUNNING") + time.sleep(0.5) + print(hq_env.command(["job", "info", "2"])) + + ts = hq_env.command(["task", "--output-mode=json", "info", "2", "0-30"], as_json=True) + print(ts) + assert ts[0]["state"] == "waiting" + + assert sum(1 if t["state"] == "running" else 0 for t in ts) == 6 + assert len(set(t["worker"] for t in ts if t["state"] == "running")) == 3