diff --git a/CHANGELOG.md b/CHANGELOG.md index e46c95abb..c43671a9e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,9 @@ * The server scheduler now contains a safety limit for computation, configurable via `--scheduler-time-limit` (default: 5s) * Better scheduling policy (prefill) for heterogenous clusters +* Workers held for a waiting high-priority task are now chosen per priority threshold, so fewer + workers are held back and more lower-priority work can run elsewhere +* More precise computation of gap for low-priority tasks ### Fixes @@ -13,6 +16,17 @@ * Fixed server crash in a specific situation when an unschedulable high-priority task occurs * Fixed server crash caused by invalid handling of prefill * Fixed canceling of prefilled tasks +* Fixed multi-node tasks that could be placed on workers from two different worker groups + when worker ids of the groups interleaved +* A held worker is no longer switched to a different worker only because small tasks were placed + into its unusable remainder; the scheduler moves a hold only to a worker that is at least as free + for the waiting task, so held capacity is not repeatedly given away +* Reservations: reservations respect priorities, and they no longer block resources + the waiting task could never use +* Fixed wrong resource amounts in this computation for tasks with resource variants on workers + that lack some resource +* Fixed constraining gap that could lead to priority inversion + ## v0.26.2 diff --git a/crates/tako/src/internal/scheduler/gap.rs b/crates/tako/src/internal/scheduler/gap.rs index 063e233f3..3cf65fcf0 100644 --- a/crates/tako/src/internal/scheduler/gap.rs +++ b/crates/tako/src/internal/scheduler/gap.rs @@ -1,7 +1,9 @@ use crate::internal::common::resources::ResourceId; use crate::internal::server::workerload::WorkerResources; use crate::internal::solver::{ConstraintType, LpSolver}; -use crate::resources::{ResourceAmount, ResourceRequestVariants, ResourceRqId, ResourceRqMap}; +use crate::resources::{ + ResourceAmount, ResourceRequest, ResourceRequestVariants, ResourceRqId, ResourceRqMap, +}; use crate::{Map, ResourceVariantId}; use hashbrown::Equivalent; use std::cell::RefCell; @@ -35,6 +37,7 @@ struct GapCacheInner { } impl GapCache { + #[cfg(test)] pub fn get_gap( &self, high_priority_rq: ResourceRqId, @@ -43,17 +46,50 @@ impl GapCache { assigned_tasks: impl Iterator, resource_rq_map: &ResourceRqMap, ) -> u32 { - let h_rqv = resource_rq_map.get(high_priority_rq); - if h_rqv.is_multi_node() { - return 0; - } let l_rqv = resource_rq_map.get(low_priority_rq); if l_rqv.is_multi_node() { return 0; } + let Some(free) = + self.gap_resources(high_priority_rq, resources, assigned_tasks, resource_rq_map) + else { + return 0; + }; + l_rqv + .requests() + .iter() + .map(|rq| free.task_max_count_for_request(rq)) + .min() + .unwrap_or(0) + } + + /// The gap allowance of a worker for a blocker: resources that tasks of any lower request may + /// jointly consume there without ever denying the blocker (`get_gap` counts how many tasks of + /// one request fit into it). `None` when no gap exists: multi-node blockers, and blockers that + /// ask for all of a resource. + pub fn gap_resources( + &self, + high_priority_rq: ResourceRqId, + resources: &WorkerResources, + assigned_tasks: impl Iterator, + resource_rq_map: &ResourceRqMap, + ) -> Option { + let h_rqv = resource_rq_map.get(high_priority_rq); + if h_rqv.is_multi_node() { + return None; + } + // Callers ask only for workers that can run the blocker; otherwise the gap is meaningless + // and the blocker may name a resource the worker lacks (out of bounds in `remove_multiple`). + debug_assert!( + h_rqv + .requests() + .iter() + .any(|rq| resources.is_capable_to_run_request(rq)), + "gap computed for a worker that cannot run the blocker" + ); let mut free: WorkerResources = if let Some(h_rq) = h_rqv.trivial_request() { if h_rq.entries().iter().any(|r| r.request.amount_is_all()) { - return 0; + return None; } let count = resources.task_max_count_for_request(h_rq); let mut resources = resources.clone(); @@ -78,19 +114,129 @@ impl GapCache { free } }; - for (rq_id, rv_id) in assigned_tasks { - if rq_id != high_priority_rq { - let rq = resource_rq_map.get(rq_id).get(rv_id); - free.remove(rq); + let occupants: Vec<&ResourceRequest> = assigned_tasks + .filter(|(rq_id, _)| *rq_id != high_priority_rq) + .map(|(rq_id, rv_id)| resource_rq_map.get(rq_id).get(rv_id)) + .collect(); + for rq in &occupants { + free.remove(rq); + } + if let Some((resource_id, exact)) = exact_single_resource_gap(h_rqv, resources, &occupants) + { + free.set(resource_id, exact); + } else if let Some(exact) = exact_gap_by_enumeration(h_rqv, resources, &occupants) { + free = exact; + } + Some(free) + } +} + +const MAX_EXACT_GAP_UNITS: u32 = 4096; + +fn exact_single_resource_gap( + h_rqv: &ResourceRequestVariants, + resources: &WorkerResources, + occupants: &[&ResourceRequest], +) -> Option<(ResourceId, ResourceAmount)> { + let h_rq = h_rqv.trivial_request()?; + let entries = h_rq.entries(); + if entries.len() != 1 { + return None; + } + let entry = &entries[0]; + let blocker = entry.request.amount_or_none_if_all()?; + let resource_id = entry.resource_id; + let capacity = resources.get(resource_id); + if blocker.fractions() != 0 || capacity.fractions() != 0 || blocker.units() == 0 { + return None; + } + let r = blocker.units(); + if r > MAX_EXACT_GAP_UNITS { + return None; + } + let capacity = capacity.units(); + + let mut reachable = vec![false; r as usize]; + reachable[0] = true; + for rq in occupants { + let amount = rq.get_amount(resource_id).unwrap_or(ResourceAmount::ZERO); + if amount.fractions() != 0 { + return None; + } + let step = (amount.units() % r) as usize; + if step == 0 { + continue; + } + let previous = reachable.clone(); + for (value, _) in previous.iter().enumerate().filter(|(_, hit)| **hit) { + reachable[(value + step) % r as usize] = true; + } + } + + let gap = (0..r) + .filter(|v| reachable[*v as usize]) + .map(|v| (capacity + r - (v % r)) % r) + .min() + .unwrap_or(0); + Some((resource_id, ResourceAmount::new(gap, 0))) +} + +const MAX_GAP_SUB_OCCUPANCIES: u32 = 2048; + +fn exact_gap_by_enumeration( + h_rqv: &ResourceRequestVariants, + resources: &WorkerResources, + occupants: &[&ResourceRequest], +) -> Option { + let h_rq = h_rqv.trivial_request()?; + if h_rq.entries().iter().any(|e| e.request.amount_is_all()) { + return None; + } + let mut groups: Map<&ResourceRequest, u32> = Map::default(); + for rq in occupants { + *groups.entry(*rq).or_default() += 1; + } + let groups = groups; + let size = groups.values().try_fold(1u32, |acc, count| { + acc.checked_mul(*count + 1) + .filter(|size| *size <= MAX_GAP_SUB_OCCUPANCIES) + })?; + let mut gap: Option = None; + let mut remaining_counts: Vec = groups.values().copied().collect(); + for _ in 0..size { + let mut remaining = resources.clone(); + for ((rq, _), k) in groups.iter().zip(&remaining_counts) { + if *k > 0 { + remaining.remove_multiple(rq, *k); } } - l_rqv - .requests() - .iter() - .map(|rq| free.task_max_count_for_request(rq)) - .min() - .unwrap_or(0) + let fit = remaining.task_max_count_for_request(h_rq); + remaining.remove_multiple(h_rq, fit); + if let Some(g) = &mut gap { + for (resource_id, amount) in remaining.iter_all_pairs() { + if amount < g.get(resource_id) { + g.set(resource_id, amount); + } + } + } else { + if h_rq + .entries() + .iter() + .all(|e| remaining.get(e.resource_id).is_zero()) + { + return Some(remaining); + } + gap = Some(remaining); + } + for ((_, count), k) in groups.iter().zip(remaining_counts.iter_mut()) { + if *k > 0 { + *k -= 1; + break; + } + *k = *count; + } } + gap } fn compute_gap_resources( @@ -107,8 +253,11 @@ fn compute_gap_resources( }; let n_resources = n_unresources + 1; let gap_res: Vec = resources - .iter_pairs() + .iter_all_pairs() .map(|(r_id, r_amount)| { + if r_amount.is_zero() { + return ResourceAmount::ZERO; + } let mut solver = LpSolver::new(false); let mut cst = vec![Vec::new(); n_resources]; let vars: Vec<_> = rqv @@ -149,6 +298,7 @@ fn compute_gap_resources( #[cfg(test)] mod tests { use crate::internal::server::core::CoreSplitMut; + use crate::resources::{CPU_RESOURCE_ID, ResourceAmount}; use std::iter; use crate::tests::utils::env::TestEnv; @@ -172,6 +322,349 @@ mod tests { .get_gap(h_rq, l_rq, &res, iter::empty(), request_map) } + fn compute_gap_occupied( + rt: &mut TestEnv, + high_task: TaskId, + low_task: TaskId, + w: WorkerId, + occupants: &[TaskId], + ) -> u32 { + let CoreSplitMut { + task_map, + worker_map, + scheduler_state, + request_map, + .. + } = rt.core().split_mut(); + let h_rq = task_map.get_task(high_task).resource_rq_id; + let l_rq = task_map.get_task(low_task).resource_rq_id; + let assigned: Vec<_> = occupants + .iter() + .map(|t| { + let task = task_map.get_task(*t); + (task.resource_rq_id, crate::ResourceVariantId::new(0)) + }) + .collect(); + let res = &worker_map.get_worker(w).resources; + scheduler_state + .gap_cache + .get_gap(h_rq, l_rq, res, assigned.into_iter(), request_map) + } + + #[test] + fn test_gap_is_one() { + let mut rt = TestEnv::new(); + let w = rt.new_worker(&WorkerBuilder::new(12)); + let r_h = rt.new_task_cpus(8); + let r_l = rt.new_task_cpus(4); + assert_eq!(compute_gap(&mut rt, r_h, r_l, w), 1, "empty worker"); + + let occupant = rt.new_task_cpus(2); + assert_eq!( + compute_gap_occupied(&mut rt, r_h, r_l, w, &[occupant]), + 0, + "two cpus occupied" + ); + } + + #[test] + fn test_gap_is_eight() { + let mut rt = TestEnv::new(); + let gpu = rt.new_named_resource("gpus"); + let w = rt.new_worker(&WorkerBuilder::new(12).res_sum("gpus", 4)); + let r_h = rt.new_task(&TaskBuilder::new().cpus(1).add_resource(gpu, 1)); + let r_l = rt.new_task_cpus(1); + assert_eq!(compute_gap(&mut rt, r_h, r_l, w), 8, "empty worker"); + + let occupant = rt.new_task_cpus(2); + assert_eq!( + compute_gap_occupied(&mut rt, r_h, r_l, w, &[occupant]), + 6, + "two cpus occupied" + ); + } + + #[test] + fn test_gap_is_safe_under_future_departures() { + let mut rt = TestEnv::new(); + let w = rt.new_worker(&WorkerBuilder::new(12)); + let unrelated = rt.new_task_cpus(3); + let r_h = rt.new_task_cpus(5); + let r_l = rt.new_task_cpus(1); + assert_eq!( + compute_gap_occupied(&mut rt, r_h, r_l, w, &[unrelated]), + 2, + "must be 12 mod 5 = 2, not the naive 9 - 5 = 4" + ); + } + + /// The case that rules out the tempting shortcut of `min(gap(C), gap(C - occupancy))`. + /// + /// C = 12, r_h = 5, occupants {1, 4}: the occupancy still present at some future point can be + /// 0, 1, 4 or 5, giving `(12 - x) mod 5` of 2, 1, 3, 2. The binding term is the *intermediate* + /// x = 1, so the gap is 1 — both extremes give 2, and granting 2 would over-commit the worker. + #[test] + fn test_gap_binding_term_can_be_an_intermediate_sub_occupancy() { + let mut rt = TestEnv::new(); + let w = rt.new_worker(&WorkerBuilder::new(12)); + let occ_a = rt.new_task_cpus(1); + let occ_b = rt.new_task_cpus(4); + let r_h = rt.new_task_cpus(5); + let r_l = rt.new_task_cpus(1); + assert_eq!( + compute_gap_occupied(&mut rt, r_h, r_l, w, &[occ_a, occ_b]), + 1, + "intermediate sub-occupancy must bind; 2 would over-grant" + ); + } + + /// Cross-check the residue DP against brute-force enumeration of every sub-occupancy, over a + /// spread of capacities, blocker widths and occupant multisets. The result must equal the + /// exact minimum -- and in particular must never exceed it, which is the unsafe direction. + #[test] + fn test_gap_matches_brute_force_over_sub_occupancies() { + for capacity in [6u32, 8, 11, 12, 16] { + for blocker in [2u32, 3, 5, 6, 7] { + if blocker > capacity { + // The gap is only computed for workers that can run the blocker + continue; + } + for occupants in [ + vec![], + vec![1u32], + vec![3], + vec![1, 4], + vec![2, 2], + vec![1, 2, 4], + vec![5, 3, 1], + ] { + if occupants.iter().sum::() > capacity { + continue; + } + // min over subsets that may remain occupied of (C - x) mod r + let mut expected = u32::MAX; + for mask in 0..(1u32 << occupants.len()) { + let x: u32 = occupants + .iter() + .enumerate() + .filter(|(i, _)| mask & (1 << i) != 0) + .map(|(_, a)| *a) + .sum(); + expected = expected.min((capacity - x) % blocker); + } + + let mut rt = TestEnv::new(); + let w = rt.new_worker(&WorkerBuilder::new(capacity)); + let occ: Vec<_> = occupants.iter().map(|a| rt.new_task_cpus(*a)).collect(); + let r_h = rt.new_task_cpus(blocker); + let r_l = rt.new_task_cpus(1); + let got = compute_gap_occupied(&mut rt, r_h, r_l, w, &occ); + assert_eq!( + got, expected, + "C={capacity} r_h={blocker} occupants={occupants:?}" + ); + } + } + } + } + + /// A blocker spanning two resources: the residue table does not apply, enumeration does. + #[test] + fn test_multi_resource_blocker_gap_is_exact() { + let mut rt = TestEnv::new(); + let gpu = rt.new_named_resource("gpus"); + let w = rt.new_worker(&WorkerBuilder::new(12).res_sum("gpus", 4)); + let r_h = rt.new_task(&TaskBuilder::new().cpus(5).add_resource(gpu, 1)); + let r_l = rt.new_task_cpus(1); + let occupant = rt.new_task_cpus(3); + // Occupant running: 9 free cpus fit one blocker, 4 cpus it cannot use. Occupant finished: + // 12 cpus fit two blockers, 2 it cannot use. Exact gap 2; the bound 12 - 10 - 3 gives 0. + assert_eq!(compute_gap_occupied(&mut rt, r_h, r_l, w, &[occupant]), 2); + } + + /// A blocker asking for a large amount of one resource is too big for the residue table, but + /// with only two occupants there are four sub-occupancies, so the gap is still exact. This is + /// the intermediate-sub-occupancy case scaled by 1000: exact 1000, where the conservative bound + /// `C - M(C) - o` would give 0. + #[test] + fn test_large_single_resource_blocker_is_exact_by_enumeration() { + let mut rt = TestEnv::new(); + let w = rt.new_worker(&WorkerBuilder::new(12_000)); + let occ_a = rt.new_task_cpus(1_000); + let occ_b = rt.new_task_cpus(4_000); + let r_h = rt.new_task_cpus(5_000); + let r_l = rt.new_task_cpus(1); + assert_eq!( + compute_gap_occupied(&mut rt, r_h, r_l, w, &[occ_a, occ_b]), + 1_000 + ); + } + + /// Too large for the residue table and too many distinct sub-occupancies to enumerate: the + /// conservative bound is used. The exact gap here would be 922; the bound is 0. + #[test] + fn test_gap_falls_back_to_bound_when_both_exact_methods_are_too_large() { + let mut rt = TestEnv::new(); + let w = rt.new_worker(&WorkerBuilder::new(12_000)); + let mut occupants = vec![rt.new_task_cpus(1_000), rt.new_task_cpus(4_000)]; + // Twelve more distinct shapes: 2^14 sub-occupancies in total. + occupants.extend((1..=12).map(|c| rt.new_task_cpus(c))); + let r_h = rt.new_task_cpus(5_000); + let r_l = rt.new_task_cpus(1); + assert_eq!(compute_gap_occupied(&mut rt, r_h, r_l, w, &occupants), 0); + } + + /// The per-resource minimum of `C - x - fit(C - x) * r_h` over every subset of occupants that + /// may still be running, for a blocker of `blocker.0` cpus + `blocker.1` gpus. Occupants are + /// `(cpus, gpus)`. Returns the allowance `(cpus, gpus)`. + fn brute_force_gap( + capacity: (u32, u32), + blocker: (u32, u32), + occupants: &[(u32, u32)], + ) -> (u32, u32) { + let mut gap = (u32::MAX, u32::MAX); + for mask in 0..(1u32 << occupants.len()) { + let (xc, xg) = occupants + .iter() + .enumerate() + .filter(|(i, _)| mask & (1 << i) != 0) + .fold((0, 0), |(c, g), (_, (oc, og))| (c + oc, g + og)); + let (fc, fg) = (capacity.0 - xc, capacity.1 - xg); + let fit = (fc / blocker.0).min(fg / blocker.1); + gap.0 = gap.0.min(fc - fit * blocker.0); + gap.1 = gap.1.min(fg - fit * blocker.1); + } + gap + } + + /// Cross-check the exact multi-resource computation against brute force over every + /// sub-occupancy, for two-resource blockers. A 1-cpu request reads off the cpu allowance, and + /// a 1-cpu + 1-gpu request the smaller of the two. + #[test] + fn test_multi_resource_gap_matches_brute_force() { + for capacity in [(8u32, 2u32), (12, 4), (16, 4), (20, 4)] { + for blocker in [(4u32, 1u32), (3, 1), (2, 2), (5, 1)] { + for occupants in [ + vec![], + vec![(2u32, 1u32)], + vec![(3, 0)], + vec![(1, 1), (2, 0)], + vec![(2, 1), (2, 1), (1, 0)], + vec![(3, 1), (1, 0), (2, 0)], + ] { + let used = occupants + .iter() + .fold((0, 0), |(c, g), (oc, og)| (c + oc, g + og)); + if used.0 > capacity.0 || used.1 > capacity.1 { + continue; + } + let expected = brute_force_gap(capacity, blocker, &occupants); + + let mut rt = TestEnv::new(); + let gpu = rt.new_named_resource("gpus"); + let w = + rt.new_worker(&WorkerBuilder::new(capacity.0).res_sum("gpus", capacity.1)); + let occ: Vec<_> = occupants + .iter() + .map(|(c, g)| { + let tb = TaskBuilder::new().cpus(*c); + rt.new_task(&if *g > 0 { tb.add_resource(gpu, *g) } else { tb }) + }) + .collect(); + let r_h = rt.new_task( + &TaskBuilder::new() + .cpus(blocker.0) + .add_resource(gpu, blocker.1), + ); + let cpu_only = rt.new_task_cpus(1); + let cpu_and_gpu = rt.new_task(&TaskBuilder::new().cpus(1).add_resource(gpu, 1)); + let context = format!("C={capacity:?} r_h={blocker:?} occupants={occupants:?}"); + assert_eq!( + compute_gap_occupied(&mut rt, r_h, cpu_only, w, &occ), + expected.0, + "cpu allowance, {context}" + ); + assert_eq!( + compute_gap_occupied(&mut rt, r_h, cpu_and_gpu, w, &occ), + expected.0.min(expected.1), + "cpu+gpu allowance, {context}" + ); + } + } + } + } + + /// Where the conservative bound is loose. A 4-cpu + 1-gpu blocker on 20 cpus and 4 gpus can + /// never use more than 16 cpus, and while a 2-cpu + 1-gpu occupant runs it gets only 3 gpus, + /// hence 12 cpus. So 4 cpus are never usable by it, whatever finishes first. The bound + /// `C - M(C) - o = 20 - 16 - 2` charges the occupant's cpus as if they stayed forever and + /// gives only 2. + #[test] + fn test_multi_resource_gap_is_larger_than_the_bound() { + let mut rt = TestEnv::new(); + let gpu = rt.new_named_resource("gpus"); + let w = rt.new_worker(&WorkerBuilder::new(20).res_sum("gpus", 4)); + let occupant = rt.new_task(&TaskBuilder::new().cpus(2).add_resource(gpu, 1)); + let r_h = rt.new_task(&TaskBuilder::new().cpus(4).add_resource(gpu, 1)); + let r_l = rt.new_task_cpus(1); + assert_eq!(compute_gap_occupied(&mut rt, r_h, r_l, w, &[occupant]), 4); + } + + /// Too many distinct sub-occupancies for a multi-resource blocker: the conservative bound is + /// used. Four 1-cpu + 1-gpu occupants and twelve cpu-only occupants of distinct sizes give + /// 5 * 2^12 sub-occupancies. The exact gap would be 200 - 78 - 16 = 106; the bound is + /// 200 - 16 - 82 = 102. + #[test] + fn test_multi_resource_gap_falls_back_to_bound_above_limit() { + let mut rt = TestEnv::new(); + let gpu = rt.new_named_resource("gpus"); + let w = rt.new_worker(&WorkerBuilder::new(200).res_sum("gpus", 4)); + let mut occupants: Vec<_> = (0..4) + .map(|_| rt.new_task(&TaskBuilder::new().cpus(1).add_resource(gpu, 1))) + .collect(); + occupants.extend((1..=12).map(|c| rt.new_task_cpus(c))); + let r_h = rt.new_task(&TaskBuilder::new().cpus(4).add_resource(gpu, 1)); + let r_l = rt.new_task_cpus(1); + assert_eq!(compute_gap_occupied(&mut rt, r_h, r_l, w, &occupants), 102); + } + + /// Multi-variant gaps are stored per resource id. A worker that lacks a resource with a lower + /// id than one it has (here no "mem", id 1, but "gpus", id 2) must still get each value at the + /// right id. + /// + /// Blocker variants: 8 cpus, or 2 cpus + 2 gpus. On 8 cpus + 3 gpus any mix uses at most + /// 8 cpus and 2 gpus, so the gap is 0 cpus and 1 gpu, and nothing for the missing "mem". + #[test] + fn test_multi_variant_gap_keeps_resource_ids_when_worker_lacks_a_resource() { + let mut rt = TestEnv::new(); + let mem = rt.new_named_resource("mem"); + let gpus = rt.new_named_resource("gpus"); + let w = rt.new_worker(&WorkerBuilder::new(8).res_sum("gpus", 3)); + let r_h = rt.new_task( + &TaskBuilder::new() + .cpus(8) + .next_variant() + .cpus(2) + .add_resource(gpus, 2), + ); + let CoreSplitMut { + task_map, + worker_map, + scheduler_state, + request_map, + .. + } = rt.core().split_mut(); + let h_rq = task_map.get_task(r_h).resource_rq_id; + let res = &worker_map.get_worker(w).resources; + let gap = scheduler_state + .gap_cache + .gap_resources(h_rq, res, iter::empty(), request_map) + .unwrap(); + assert_eq!(gap.get(CPU_RESOURCE_ID), ResourceAmount::ZERO, "cpus"); + assert_eq!(gap.get(mem), ResourceAmount::ZERO, "mem"); + assert_eq!(gap.get(gpus), ResourceAmount::new_units(1), "gpus"); + } + #[test] fn test_compute_gap() { let mut rt = TestEnv::new(); diff --git a/crates/tako/src/internal/scheduler/solver.rs b/crates/tako/src/internal/scheduler/solver.rs index 91abb24f6..b5ff829be 100644 --- a/crates/tako/src/internal/scheduler/solver.rs +++ b/crates/tako/src/internal/scheduler/solver.rs @@ -1,10 +1,14 @@ -use crate::internal::common::resources::{ResourceId, ResourceRequest, ResourceRequestVariants}; +use crate::internal::common::resources::{ + ResourceAmount, ResourceId, ResourceRequest, ResourceRequestVariants, +}; use crate::internal::scheduler::TaskBatch; +use crate::internal::scheduler::state::SchedulerState; use crate::internal::server::core::{Core, CoreSplit}; +use crate::internal::server::taskmap::TaskMap; 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::resources::{CPU_RESOURCE_ID, ResourceRqId, ResourceRqMap}; use crate::{Map, ResourceVariantId, Set, WorkerId}; use thin_vec::ThinVec; @@ -34,11 +38,6 @@ impl SchedulingSolution { } } -struct Candidates { - demand: u32, - candidates: Vec<(f64, WorkerId)>, -} - pub(crate) fn run_scheduling_solver( core: &Core, now: std::time::Instant, @@ -93,24 +92,38 @@ pub(crate) fn run_scheduling_solver( let mut placements: Map<(WorkerId, ResourceRqId, ResourceVariantId), Variable> = Map::new(); let mut tasks_count_vars: Map> = Map::new(); + // Reservation variables by (worker, request). A reservation stands for one task of its request + // set aside on that worker, so it also enters the request's own priority conditions there. + let mut reservations: Map<(WorkerId, ResourceRqId), Variable> = Map::new(); 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(), - }, - ); - } + // Held workers are chosen before any variable exists, because a reservation variable is only + // created on a held worker. A reservation elsewhere would count toward its blocker without + // holding anything back: on a worker with no free resources it costs nothing but its small + // penalty, and with no worker held (the blocker fits somewhere right now) it would let a + // lower-priority request take the very capacity the blocker was counted as using. + let reserved = held_workers( + &workers, + task_batches, + request_map, + task_map, + scheduler_state, + now, + ); + // The reservation penalty is expressed relative to the placements a reservation can unlock. + // Placement weights shrink as the cluster's free resources grow, so a fixed penalty would + // outweigh every placement on a large cluster and no reservation would ever pay for itself. + let reservation_scale = + placement_weight_lower_bound(&workers, task_batches, request_map, &resource_sums); + // Create worker-task placements + let mut worker_reservations: Vec = Vec::new(); for (w_idx, worker) in workers.iter().enumerate() { worker_cpu_constraint_no_reserves.clear(); + worker_reservations.clear(); for batch in task_batches.iter() { let rqv = request_map.get(batch.resource_rq_id); let mut has_variant = false; @@ -133,14 +146,11 @@ pub(crate) fn run_scheduling_solver( ); placements.insert((worker.id, batch.resource_rq_id, v_idx), v); // Insert into worker resource constraints - for (r, amount) in worker.resources.iter_pairs() { + for (r, amount) in worker.resources.iter_nonzero_pairs() { worker_res_constraint[r.as_usize()].push((v, amount.as_f64())); } } - } else if !worker.is_request_blocked(batch.resource_rq_id, v_idx) - && worker.has_time_to_run(rq.min_time(), now) - && worker.have_immediate_resources_for_rq(rq) - { + } else if sn_variant_fits_now(worker, batch.resource_rq_id, v_idx, rq, now) { has_variant = true; set_placement_name(&mut solver, worker.id, batch.resource_rq_id, v_idx); let v = @@ -167,44 +177,55 @@ 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) + && reserved + .get(&batch.resource_rq_id) + .is_some_and(|held| held.order.contains(&worker.id)) && let Some(a) = worker.sn_assignment() { - let weight = -((n_workers - w_idx) as f64 / (n_workers * 1024) as f64); + let weight = + -reservation_scale * (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); + reservations.insert((worker.id, batch.resource_rq_id), v); + worker_reservations.push(v); tasks_count_vars .entry(batch.resource_rq_id) .or_default() .push(v); - for (res_id, count) in a.free_resources.iter_pairs() { - worker_res_constraint[res_id.as_usize()].push((v, count.as_f64())); + // A reservation withholds the free resources the blocker could use, but not its gap + // allowance: that capacity is lost to the blocker anyway, and the priority + // conditions already let lower-priority gap tasks use it. + let gap = worker_gap( + worker, + batch.resource_rq_id, + request_map, + task_map, + scheduler_state, + ); + for (res_id, count) in a.free_resources.iter_nonzero_pairs() { + let withheld = gap + .as_ref() + .map_or(count, |g| count.saturating_sub(g.get(res_id))); + if !withheld.is_zero() { + worker_res_constraint[res_id.as_usize()].push((v, withheld.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)); } } + // A worker can later host only one of the blockers it is reserved for, so it may serve + // at most one. The resource constraints do not ensure this: on a worker with no free + // resources a reservation consumes nothing, and any number of them would fit. + if worker_reservations.len() > 1 { + solver.set_name(|| format!("w{}: at most one reservation", worker.id)); + solver.add_constraint( + ConstraintType::Max, + 1.0, + worker_reservations.iter().map(|v| (*v, 1.0)), + ); + } + if worker.configuration.min_utilization > 0.001 { add_min_utilization(&mut solver, worker, &mut worker_cpu_constraint_no_reserves); } @@ -226,30 +247,6 @@ 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(); @@ -313,9 +310,13 @@ pub(crate) fn run_scheduling_solver( Some(new_v) }; - let mut zero_cond = Vec::new(); - let mut zero_cond_reserved = Vec::new(); + // Terms of a condition's left-hand side: every placement, minus what sits in a worker's gap. + let mut cond_terms: Vec<(Variable, f64)> = Vec::new(); + let mut cond_terms_reserved: Vec<(Variable, f64)> = Vec::new(); let mut blocked_by_unbounded: Set = Set::new(); + // Gap parts per (worker, blocker, resource): the gap allowance and the terms that share it. + let mut shared_gap: Map)> = Map::new(); + let mut shared_gap_seen: Set<(WorkerId, ResourceRqId, ResourceRqId)> = Set::new(); for batch in task_batches.iter() { let Some(task_counts) = tasks_count_vars.get(&batch.resource_rq_id) else { @@ -335,8 +336,8 @@ pub(crate) fn run_scheduling_solver( blocked_by_unbounded.clear(); for cut in &batch.cuts { for (blocker_rq_id, blocking_size) in &cut.blockers { - zero_cond.clear(); - zero_cond_reserved.clear(); + cond_terms.clear(); + cond_terms_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() { @@ -344,7 +345,7 @@ pub(crate) fn run_scheduling_solver( task_counts_per_group.get(&(batch.resource_rq_id, group_name.as_str())) && group.is_capable_to_run(blocker_rqv, now, worker_map) { - zero_cond.push(*v); + cond_terms.push((*v, 1.0)); } } } else { @@ -355,9 +356,8 @@ pub(crate) fn run_scheduling_solver( if !w.is_capable_to_run_rqv(blocker_rqv, now) { continue; } - let gap = scheduler_state.gap_cache.get_gap( + let gap_resources = scheduler_state.gap_cache.gap_resources( *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); @@ -368,106 +368,114 @@ pub(crate) fn run_scheduling_solver( }), request_map, ); + let gap = gap_resources + .as_ref() + .map(|g| gap_count(g, batch_rqv)) + .unwrap_or(0); // 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(|| { - format!( - "w{}: if #rq{blocker_rq_id} < {s} then limit #rq{} to {} + {} (gap) where both rqs may run", - w.id, batch.resource_rq_id, cut.size, gap - ) - }); - constraint_extra_var( - &mut solver, - ConstraintType::Max, - cut_size + batch_size + gap as f64, - vars, - blocking_v, - batch_size, - ); - } 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}", - w.id, batch.resource_rq_id, gap, cut.size, - ) - }); - solver.add_constraint( - ConstraintType::Max, - cut_size + gap as f64, - vars.into_iter().map(|v| (v, 1.0)), - ); - } - } else { - for v_id in batch_rqv.variant_ids() { - if let Some(var) = - placements.get(&(w.id, batch.resource_rq_id, v_id)) + let is_reserved = blocking_size.is_some_and(|threshold| { + reserved + .get(blocker_rq_id) + .is_some_and(|held| held.holds(w.id, threshold)) + }); + // A reservation counts like a placement of its request: while the blocker is + // unserved on this worker, the request may not set capacity aside here any + // more than it may start a task here. + let reservation = reservations.get(&(w.id, batch.resource_rq_id)).copied(); + let placement_vars = batch_rqv + .variant_ids() + .filter_map(|v_id| { + placements + .get(&(w.id, batch.resource_rq_id, v_id)) + .map(|v| (v_id, *v)) + }) + .collect::>(); + if placement_vars.is_empty() && reservation.is_none() { + continue; + } + let placed = placement_vars + .iter() + .map(|(_, v)| (*v, 1.0)) + .chain(reservation.map(|v| (v, 1.0))); + cond_terms.extend(placed.clone()); + if is_reserved { + cond_terms_reserved.extend(placed); + } + if gap == 0 { + // Nothing can sit in a gap here, so every task placed counts against + // the before-count. No variable and no constraint for this worker. + continue; + } + // Tasks of this request sitting in the worker's gap. They do not count + // against the before-count, so they are subtracted from the condition. + // At most `gap` of them fit, which is a bound of the variable itself. + solver.set_name(|| { + format!( + "w{}: #rq{} within the {gap} gap of rq{blocker_rq_id}", + w.id, batch.resource_rq_id + ) + }); + let gap_var = solver.add_variable(0.0, 0.0, gap as f64); + // Only tasks that really run here may sit in the gap. Without this, a + // request could claim gap usage it does not have, and so shrink the + // condition below and the shared gap constraint for other requests. + solver.set_name(|| { + format!( + "w{}: #rq{} in the gap is at most what runs here", + w.id, batch.resource_rq_id + ) + }); + constraint_extra_var( + &mut solver, + ConstraintType::Min, + 0.0, + placement_vars.iter().map(|(_, v)| *v).chain(reservation), + gap_var, + -1.0, + ); + cond_terms.push((gap_var, -1.0)); + if is_reserved { + cond_terms_reserved.push((gap_var, -1.0)); + } + // All requests blocked by this blocker share the worker's gap, so their + // gap parts are collected per resource. Recorded once per (worker, + // blocker, request): the gap does not depend on the cut. + if let Some(gap_resources) = gap_resources.as_ref() + && shared_gap_seen.insert((w.id, *blocker_rq_id, batch.resource_rq_id)) + { + for (resource_id, gap_amount) in gap_resources.iter_all_pairs() { + if gap_amount >= sn_assignment.free_resources.get(resource_id) { + // The worker's own resource limit is at least as strict. + continue; + } + // Gap tasks are charged at the largest variant amount, so a mix + // of variants can never use more of the gap than counted. + let largest = placement_vars + .iter() + .map(|(v_id, _)| { + batch_rqv + .get(*v_id) + .get_amount(resource_id) + .unwrap_or_else(|| w.resources.get(resource_id)) + }) + .max(); + if let Some(largest) = largest + && !largest.is_zero() { - zero_cond.push(*var); - if is_reserved { - zero_cond_reserved.push(*var); - } + shared_gap + .entry((w.id, *blocker_rq_id, resource_id)) + .or_insert_with(|| (gap_amount, Vec::new())) + .1 + .push((gap_var, largest.as_f64())); } } } - /*for v_id in batch_rqv.variant_ids() { - if let Some(var) = placements.get(&(w.id, batch.resource_rq_id, v_id)) { - let gap = scheduler_cache.gap_cache.get_gap( - *blocker_rq_id, - batch.resource_rq_id, - v_id, - &w.resources, - request_map, - ); - if gap > 0 { - let cut_size = cut.size as f64; - if let Some(s) = blocking_size { - let blocking_v = get_bvar(&mut solver, *blocker_rq_id, *s); - solver.set_name(|| { - format!( - "w{}: if #rq{blocker_rq_id} < {s} then limit #rq{} to {} + {} (gap) where both rqs may run", - w.id, batch.resource_rq_id, cut.size, gap - ) - }); - solver.add_constraint( - ConstraintType::Max, - cut_size + batch_size + gap as f64, - [(*var, 1.0), (blocking_v, batch_size)].into_iter(), - ); - } else { - solver.set_name(|| { - format!( - "w{}: limit #rq{} to {} + {} (gap) where it can run with rq{blocker_rq_id}", - w.id, batch.resource_rq_id, gap, cut.size, - ) - }); - solver.add_constraint( - ConstraintType::Max, - cut_size + gap as f64, - [(*var, 1.0)].into_iter(), - ); - } - } else { - zero_cond.push(*var); - } - } - }*/ } } // 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() { + if !cond_terms_reserved.is_empty() { solver.set_name(|| { format!( "limit #rq{} to {} on workers held for rq{blocker_rq_id}", @@ -477,12 +485,14 @@ pub(crate) fn run_scheduling_solver( solver.add_constraint( ConstraintType::Max, cut.size as f64, - zero_cond_reserved.iter().map(|v| (*v, 1.0)), + cond_terms_reserved.iter().copied(), ); } - if zero_cond.is_empty() { + if cond_terms.is_empty() { continue; } + // The before-count is one number for the whole cluster: every placement counts + // against it, except what sits in the gap of its worker. if let Some(s) = blocking_size && let Some(blocking_v) = get_bvar(&mut solver, *blocker_rq_id, *s) { @@ -493,15 +503,18 @@ pub(crate) fn run_scheduling_solver( ) }); let cut_size = cut.size as f64; - constraint_extra_var( - &mut solver, + solver.add_constraint( ConstraintType::Max, batch_size + cut_size, - zero_cond.iter().copied(), - blocking_v, - batch_size, + cond_terms + .iter() + .copied() + .chain(std::iter::once((blocking_v, batch_size))), ); - } else if blocking_size.is_none() && !blocked_by_unbounded.contains(blocker_rq_id) { + } else if blocking_size.is_none() + && (cond_terms.iter().any(|(_, coef)| *coef < 0.0) + || !blocked_by_unbounded.contains(blocker_rq_id)) + { blocked_by_unbounded.insert(*blocker_rq_id); solver.set_name(|| { format!( @@ -512,13 +525,24 @@ pub(crate) fn run_scheduling_solver( solver.add_constraint( ConstraintType::Max, cut.size as f64, - zero_cond.iter().map(|v| (*v, 1.0)), + cond_terms.iter().copied(), ); } } } } + // One gap per worker and blocker, shared by every request the blocker blocks. + for ((worker_id, blocker_rq_id, resource_id), (gap_amount, terms)) in shared_gap { + solver.set_name(|| { + format!( + "w{worker_id}: gap of rq{blocker_rq_id} on resource {} is {gap_amount}", + resource_id.as_num() + ) + }); + solver.add_constraint(ConstraintType::Max, gap_amount.as_f64(), terms.into_iter()); + } + let mut result = SchedulingSolution::default(); let Some((solution, is_optimal)) = solver.solve_bounded(scheduler_state.config.mip_time_limit) else { @@ -533,22 +557,26 @@ pub(crate) fn run_scheduling_solver( let v_id = ResourceVariantId::new(0); let n_nodes = rqv.get(v_id).n_nodes() as usize; let mut ws: Vec> = Vec::new(); + // Chunks are cut per worker group: the group constraint makes each group's count a + // multiple of `n_nodes`, but worker ids of different groups may interleave, so cutting + // across all workers in id order could give a task nodes from two groups. + let mut open: Map<&str, ThinVec> = Map::new(); for worker in &workers { if let Some(v) = placements.get(&(worker.id, resource_rq_id, v_id)) { let count = solution.get_value(*v).round() as u32; if count > 0 { - if let Some(last) = ws.last_mut() - && last.len() < n_nodes - { - last.push(worker.id); - } else { - let mut workers = ThinVec::with_capacity(n_nodes); - workers.push(worker.id); - ws.push(workers); + let group = worker.configuration.group.as_str(); + let chunk = open + .entry(group) + .or_insert_with(|| ThinVec::with_capacity(n_nodes)); + chunk.push(worker.id); + if chunk.len() == n_nodes { + ws.push(open.remove(group).unwrap()); } } } } + assert!(open.is_empty()); if !ws.is_empty() { result.mn_workers.insert((resource_rq_id, v_id), ws); } @@ -629,6 +657,208 @@ fn add_min_utilization( worker_res_constraint.pop(); } +/// Whether a single-node variant of a request can be placed on the worker in this round. +/// Shared by placement creation and by `held_workers`, which must agree on it exactly. +fn sn_variant_fits_now( + worker: &Worker, + rq_id: ResourceRqId, + v_idx: ResourceVariantId, + rq: &ResourceRequest, + now: std::time::Instant, +) -> bool { + !worker.is_request_blocked(rq_id, v_idx) + && worker.has_time_to_run(rq.min_time(), now) + && worker.have_immediate_resources_for_rq(rq) +} + +/// The gap allowance of `worker` for blocker `rq_id`, given what currently runs there. +fn worker_gap( + worker: &Worker, + rq_id: ResourceRqId, + request_map: &ResourceRqMap, + task_map: &TaskMap, + scheduler_state: &SchedulerState, +) -> Option { + let a = worker.sn_assignment()?; + scheduler_state.gap_cache.gap_resources( + rq_id, + &worker.resources, + a.assigned_tasks.iter().map(|task_id| { + let t = task_map.get_task(*task_id); + ( + t.resource_rq_id, + t.assigned_placement(&scheduler_state.redirects).unwrap().1, + ) + }), + request_map, + ) +} + +/// Workers held for one single-node blocker, most empty first (ties to the higher worker id). +struct Held { + /// Blocker tasks that the workers' currently free resources can host. + placeable: u32, + /// Held workers for the blocker's largest threshold; any smaller threshold holds a prefix. + order: Vec, +} + +impl Held { + /// Whether a priority condition that asks for `threshold` blocker tasks holds this worker: + /// it needs one held worker per blocker task it counts that cannot be placed now. + fn holds(&self, worker_id: WorkerId, threshold: u32) -> bool { + let needed = threshold.saturating_sub(self.placeable) as usize; + self.order.iter().take(needed).any(|w| *w == worker_id) + } +} + +/// Workers held for each single-node blocker. +/// +/// Holding is sized by the thresholds of the priority conditions that name the blocker, not by +/// its batch: a condition asking for `s` blocker tasks needs `s - placeable` held workers, and +/// blocker tasks below every request they could block need none. Sizing by the whole batch +/// would hold capacity for blocker tasks that the blocked request outranks. +fn held_workers( + workers: &[&Worker], + task_batches: &[TaskBatch], + request_map: &ResourceRqMap, + task_map: &TaskMap, + scheduler_state: &SchedulerState, + now: std::time::Instant, +) -> Map { + let mut max_threshold: Map = Map::new(); + for batch in task_batches { + for cut in &batch.cuts { + for (blocker_rq_id, blocking_size) in &cut.blockers { + if let Some(size) = blocking_size { + let entry = max_threshold.entry(*blocker_rq_id).or_default(); + *entry = (*entry).max(*size); + } + } + } + } + + let mut held = Map::new(); + for batch in task_batches { + let rqv = request_map.get(batch.resource_rq_id); + if !batch.is_blocker || rqv.is_multi_node() { + continue; + } + let Some(threshold) = max_threshold.get(&batch.resource_rq_id).copied() else { + continue; + }; + let mut placeable = 0u32; + let mut candidates: Vec<(f64, WorkerId)> = Vec::new(); + for worker in workers { + let Some(a) = worker.sn_assignment() else { + continue; + }; + let fits_now = rqv.requests_with_ids().any(|(v_idx, rq)| { + sn_variant_fits_now(worker, batch.resource_rq_id, v_idx, rq, now) + }); + if fits_now { + placeable = placeable.saturating_add(a.free_resources.task_max_count(rqv)); + } else if worker.is_capable_to_run_rqv(rqv, now) { + // The blocker cannot run here now, but this worker could host it once it drains. + // Rank by the free resources the blocker can use: its gap is lost to it anyway. + // Tasks placed into the gap lower free resources and gap equally, so they never + // move the hold, and a finishing task never lowers the rank. A held worker + // therefore loses its hold only to a worker that offers the blocker more. + let mut usable = a.free_resources.clone(); + if let Some(gap) = worker_gap( + worker, + batch.resource_rq_id, + request_map, + task_map, + scheduler_state, + ) { + for (res_id, amount) in gap.iter_nonzero_pairs() { + usable.set(res_id, usable.get(res_id).saturating_sub(amount)); + } + } + candidates.push((fit_ratio(&usable, rqv), worker.id)); + } + } + let needed = threshold.saturating_sub(placeable) as usize; + if needed == 0 { + continue; + } + candidates.sort_unstable_by(|(fit_a, w_a), (fit_b, w_b)| { + fit_b.total_cmp(fit_a).then(w_b.cmp(w_a)) + }); + held.insert( + batch.resource_rq_id, + Held { + placeable, + order: candidates + .into_iter() + .take(needed) + .map(|(_fit, w_id)| w_id) + .collect(), + }, + ); + } + held +} + +/// A lower bound on the objective weight of any single-node placement this round. +/// +/// `create_sn_var` weighs a placement by its request's share of the free resources, the worker's +/// compaction bias `(n - i_w) / n`, and the request weight. The bias is smallest on the last +/// worker, `1 / n`, and a request for *all* of a resource is valued at the worker's total, which is +/// at least the smallest total any worker has. Only positive weights count; if there are none, +/// nothing can be unlocked and the scale is irrelevant. +fn placement_weight_lower_bound( + workers: &[&Worker], + task_batches: &[TaskBatch], + request_map: &ResourceRqMap, + resource_sums: &[f64], +) -> f64 { + let n_workers = workers.len(); + if n_workers == 0 { + return 1.0; + } + let smallest_total = |r: ResourceId| { + workers + .iter() + .map(|w| w.resources.get(r).as_f64()) + .filter(|amount| *amount > 0.0) + .fold(f64::INFINITY, f64::min) + }; + let mut bound = f64::INFINITY; + for batch in task_batches { + for rq in request_map.get(batch.resource_rq_id).requests() { + if rq.is_multi_node() { + continue; + } + let share: f64 = rq + .entries() + .iter() + .map(|e| { + let global = resource_sums[e.resource_id.as_usize()]; + if global < 0.000001 { + return 0.0; + } + let amount = e + .request + .amount_or_none_if_all() + .map(|a| a.as_f64()) + .unwrap_or_else(|| smallest_total(e.resource_id)); + if amount.is_finite() { + amount / global + } else { + 0.0 + } + }) + .sum(); + let weight = share * rq.weight().as_f64() / n_workers as f64; + if weight > 0.0 { + bound = bound.min(weight); + } + } + } + if bound.is_finite() { bound } else { 1.0 } +} + fn fit_ratio(free: &WorkerResources, rqv: &ResourceRequestVariants) -> f64 { rqv.requests() .iter() @@ -690,7 +920,7 @@ fn create_mn_var( ) -> Variable { let weight = worker .resources - .iter_pairs() + .iter_nonzero_pairs() .map(|(r, amount)| { let global = resource_sums[r.as_usize()]; if global < 0.000001 { @@ -706,6 +936,21 @@ fn create_mn_var( solver.add_bool_variable(weight) } +/// How many tasks of `rqv` fit into a worker's gap allowance for some blocker. +/// Worker, blocker request and resource: one shared gap allowance. +type SharedGapKey = (WorkerId, ResourceRqId, ResourceId); + +fn gap_count(gap_resources: &WorkerResources, rqv: &ResourceRequestVariants) -> u32 { + if rqv.is_multi_node() { + return 0; + } + rqv.requests() + .iter() + .map(|rq| gap_resources.task_max_count_for_request(rq)) + .min() + .unwrap_or(0) +} + fn constraint_extra_var( solver: &mut LpSolver, constraint_type: ConstraintType, diff --git a/crates/tako/src/internal/server/workerload.rs b/crates/tako/src/internal/server/workerload.rs index cdcbb0254..888709f03 100644 --- a/crates/tako/src/internal/server/workerload.rs +++ b/crates/tako/src/internal/server/workerload.rs @@ -30,18 +30,20 @@ impl WorkerResources { .unwrap_or(ResourceAmount::ZERO) } - pub(crate) fn iter_pairs(&self) -> impl Iterator { - self.n_resources - .iter() - .copied() + pub(crate) fn iter_all_pairs(&self) -> impl Iterator { + self.iter_amounts() .enumerate() - .filter_map(|(idx, c)| { - if !c.is_zero() { - Some((ResourceId::new(idx as u32), c)) - } else { - None - } - }) + .map(|(idx, c)| (ResourceId::new(idx as u32), c)) + } + + pub(crate) fn iter_nonzero_pairs(&self) -> impl Iterator { + self.iter_amounts().enumerate().filter_map(|(idx, c)| { + if !c.is_zero() { + Some((ResourceId::new(idx as u32), c)) + } else { + None + } + }) } pub(crate) fn iter_amounts(&self) -> impl Iterator { @@ -153,6 +155,12 @@ impl WorkerResources { .sum::() } + /// Overwrite one resource's amount. Used by the gap computation, which derives an exact + /// value for the blocker's resource and keeps the derived-by-subtraction value elsewhere. + pub fn set(&mut self, resource_id: ResourceId, amount: ResourceAmount) { + self.n_resources[resource_id] = amount; + } + pub fn remove(&mut self, rq: &ResourceRequest) { for entry in rq.entries() { if let Some(amount) = entry.request.amount_or_none_if_all() { diff --git a/crates/tako/src/internal/tests/test_scheduler_mn.rs b/crates/tako/src/internal/tests/test_scheduler_mn.rs index 3a48cf225..06cc83059 100644 --- a/crates/tako/src/internal/tests/test_scheduler_mn.rs +++ b/crates/tako/src/internal/tests/test_scheduler_mn.rs @@ -354,3 +354,30 @@ fn test_schedule_mn_and_sn4() { assert!(rt.task(t1).is_mn_running()); assert!(rt.task(t2).is_assigned()); } + +#[test] +fn test_mn_task_stays_within_one_group_when_group_ids_interleave() { + // Worker ids of two groups interleave (a, b, a, b), as when two allocations start at the + // same time. Each group can run one 2-node task, and each task must use workers of a single + // group: a task spanning two allocations would not share an interconnect. + let mut rt = TestEnv::new(); + for group in ["a", "b", "a", "b"] { + rt.new_worker(&WorkerBuilder::new(1).group(group)); + } + rt.new_task(&TaskBuilder::new().n_nodes(2)); + rt.new_task(&TaskBuilder::new().n_nodes(2)); + let solution = rt.schedule_solution(); + + let tasks: Vec<_> = solution.mn_workers.values().flatten().collect(); + assert_eq!(tasks.len(), 2); + for workers in tasks { + let groups: Vec<&str> = workers + .iter() + .map(|w| rt.worker(*w).configuration.group.as_str()) + .collect(); + assert!( + groups.iter().all(|g| *g == groups[0]), + "multi-node task {workers:?} spans worker groups {groups:?}" + ); + } +} diff --git a/crates/tako/src/internal/tests/test_scheduler_sn.rs b/crates/tako/src/internal/tests/test_scheduler_sn.rs index a23b137d5..4abce4e93 100644 --- a/crates/tako/src/internal/tests/test_scheduler_sn.rs +++ b/crates/tako/src/internal/tests/test_scheduler_sn.rs @@ -1858,6 +1858,377 @@ fn test_schedule_reservation_leaves_other_workers_for_backfill() { ); } +/// Narrow tasks assigned to each of `workers`, in order. +fn narrow_per_worker(rt: &TestEnv, narrow: &[TaskId], workers: &[WorkerId]) -> Vec { + workers + .iter() + .map(|w| { + narrow + .iter() + .filter(|t| { + matches!(rt.task(**t).state, + TaskRuntimeState::Assigned { worker_id, .. } if worker_id == *w) + }) + .count() + }) + .collect() +} + +#[test] +fn test_schedule_placeable_blocker_holds_nothing() { + // A wide blocker that a free worker can host right now needs no held worker, so capable but + // partially occupied workers keep backfilling. + let mut rt = TestEnv::new(); + let ws = rt.new_workers(4, &WorkerBuilder::new(8)); + for w in &ws[1..] { + rt.new_task_running(&TaskBuilder::new().cpus(6), *w); + } + let blocker = rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + let narrow = rt.new_tasks(30, &TaskBuilder::new().cpus(1)); + rt.schedule(); + + assert!( + matches!(rt.task(blocker).state, + TaskRuntimeState::Assigned { worker_id, .. } if worker_id == ws[0]), + "the blocker starts on the free worker" + ); + assert_eq!( + narrow_per_worker(&rt, &narrow, &ws[1..]), + vec![2, 2, 2], + "no worker is held, so every partial worker backfills its 2 free cpus" + ); +} + +#[test] +fn test_schedule_lower_priority_cannot_take_capacity_of_higher() { + // One free worker can host either the higher-priority 4-cpu task or the lower-priority 8-cpu + // task, but not both, so priority must give it to the 4-cpu task. + // + // The saturated worker is what makes this fail. It is capable of the 4-cpu task but has no + // free cpus, so the solver gets a reservation variable for that task there, and reserving a + // worker with nothing free costs no capacity. Setting it discharges the priority condition + // "8-cpu tasks may run only once the 4-cpu task is served", after which the 8-cpu task is the + // more valuable placement. Held workers do not close this: held-worker depth is the number of + // blockers no worker can host *before* the solve, and the free worker is counted as hosting + // the 4-cpu task although the solve then gives it to the 8-cpu one. + // + // Without the saturated worker the 4-cpu task is placed correctly. + let mut rt = TestEnv::new(); + let ws = rt.new_workers(2, &WorkerBuilder::new(8)); + for _ in 0..8 { + rt.new_task_running(&TaskBuilder::new().cpus(1), ws[1]); + } + let high = rt.new_task(&TaskBuilder::new().cpus(4).user_priority(11)); + let wide = rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + rt.schedule(); + + assert!( + matches!(rt.task(high).state, + TaskRuntimeState::Assigned { worker_id, .. } if worker_id == ws[0]), + "the higher-priority task must take the free worker; high = {:?}, wide = {:?}", + rt.task(high).state, + rt.task(wide).state + ); + assert!( + !rt.task(wide).is_assigned(), + "the lower-priority task must wait; wide = {:?}", + rt.task(wide).state + ); +} + +#[test] +fn test_schedule_reservation_cannot_take_capacity_of_higher_blocker() { + // A reservation for a request must obey that request's own priority conditions, exactly as + // its placements do. Here `x` (8 cpus, priority 10) cannot run anywhere now, so the emptiest + // capable worker `w1` is held for it. `h` (4 cpus + a "foo" only `w1` has, priority 20) + // outranks `x` and fits into `w1`'s free cpus. Reserving `w1` for `x` would consume those + // cpus and leave `h` waiting behind a lower-priority request. + // + // The objective alone does not prevent it. `h` is incapable of the other workers, so a + // reservation for `x` releases their freed cpus for narrow work without `h` blocking it, and + // eight workers' worth of backfill outweighs placing `h`. + let mut rt = TestEnv::new(); + rt.new_named_resource("foo"); + let w1 = rt.new_worker(&WorkerBuilder::new(8).res_sum("foo", 1000)); + rt.new_task_running(&TaskBuilder::new().cpus(4), w1); + let others = rt.new_workers(8, &WorkerBuilder::new(8)); + for w in &others { + rt.new_task_running(&TaskBuilder::new().cpus(6), *w); + } + let h = rt.new_task( + &TaskBuilder::new() + .cpus(4) + .add_resource(1, 1) + .user_priority(20), + ); + let x = rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + rt.new_tasks(40, &TaskBuilder::new().cpus(1)); + rt.schedule(); + + assert!( + matches!(rt.task(h).state, + TaskRuntimeState::Assigned { worker_id, .. } if worker_id == w1), + "the higher-priority task must take the free cpus on w1; h = {:?}, x = {:?}", + rt.task(h).state, + rt.task(x).state + ); +} + +#[test] +fn test_schedule_held_worker_not_refilled_when_one_reservation_serves_blocker() { + // Two 8-cpu tasks share one request, but the narrow request's priority lies between them: + // `x_high` (10) > narrow (5) > `x_low` (3). The held count covers the whole batch, so two + // workers are held, while the narrow request is blocked only by `x_high`, i.e. at threshold 1. + // One reservation therefore serves the blocker for the narrow request, and a second held + // worker is left with free capacity. + // + // `w0` is the emptiest worker and should accumulate for `x_high`. If narrow work may enter a + // held worker whenever the blocker is served elsewhere, the solver reserves `w1` instead (a + // reservation on a higher index is cheaper) and refills `w0`, so `x_high` waits for the fuller + // worker to drain. + let mut rt = TestEnv::new(); + let ws = rt.new_workers(2, &WorkerBuilder::new(8)); + let mut running: Vec> = [4, 6] + .iter() + .zip(&ws) + .map(|(n, w)| { + (0..*n) + .map(|_| rt.new_task_running(&TaskBuilder::new().cpus(1), *w)) + .collect() + }) + .collect(); + let x_high = rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + rt.new_tasks(40, &TaskBuilder::new().cpus(1).user_priority(5)); + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(3)); + + // `w0` needs 4 finishes to fit an 8-cpu task, `w1` needs 6. + for _tick in 0..=4 { + rt.schedule(); + if rt.task(x_high).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!( + "the priority-10 task did not start once the emptiest worker drained ({})", + free.join(", ") + ); +} + +#[test] +fn test_schedule_no_worker_held_for_blocker_tasks_the_blocked_request_outranks() { + // `x_high` (10) > narrow (5) > `x_low` (3), both x tasks 8 cpus and unplaceable now. Narrow is + // blocked only by `x_high`, so only one worker is needed for it: the emptiest, `w0`. Holding a + // second worker for `x_low` would keep narrow off it although narrow outranks `x_low`. + let mut rt = TestEnv::new(); + let ws = rt.new_workers(2, &WorkerBuilder::new(8)); + for (n, w) in [4, 6].iter().zip(&ws) { + for _ in 0..*n { + rt.new_task_running(&TaskBuilder::new().cpus(1), *w); + } + } + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + let narrow = rt.new_tasks(40, &TaskBuilder::new().cpus(1).user_priority(5)); + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(3)); + rt.schedule(); + + assert_eq!( + narrow_per_worker(&rt, &narrow, &ws), + vec![0, 2], + "w0 is held for x_high; w1 must be free for narrow work" + ); +} + +#[test] +fn test_schedule_held_worker_not_refilled_at_a_shallower_threshold() { + // Two lower requests block on the same 8-cpu request at different thresholds: + // `x_a` (10) > `l1` (9) > `x_b` (8) > `l2` (7). `l2` needs both x tasks served, so two workers + // are held and both carry a reservation variable; `l1` needs only `x_a`, i.e. one. A single + // reservation on the fuller `w1` serves `x_a` for `l1`, so the emptiest `w0` must stay closed to + // `l1` by the held-worker constraint itself, or `l1` refills it and `x_a` waits for `w1`. + let mut rt = TestEnv::new(); + let ws = rt.new_workers(2, &WorkerBuilder::new(8)); + let mut running: Vec> = [4, 6] + .iter() + .zip(&ws) + .map(|(n, w)| { + (0..*n) + .map(|_| rt.new_task_running(&TaskBuilder::new().cpus(1), *w)) + .collect() + }) + .collect(); + let x_a = rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + rt.new_tasks(40, &TaskBuilder::new().cpus(1).user_priority(9)); + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(8)); + rt.new_tasks(40, &TaskBuilder::new().cpus(2).user_priority(7)); + + for _tick in 0..=4 { + rt.schedule(); + if rt.task(x_a).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!( + "the priority-10 task did not start once the emptiest worker drained ({})", + free.join(", ") + ); +} + +#[test] +fn test_schedule_reservation_pays_on_a_large_cluster() { + // Placement weights scale with 1 / (free cpus in the cluster), so on a large cluster each + // placement is worth very little. The reservation penalty must scale with them: a reservation + // that unlocks placements has to pay for itself regardless of cluster size. + // + // `x` (8 cpus + 1 "foo", priority 10) fits nowhere now. `w0` is held for it. `w1` is capable + // but not held, so narrow work there is blocked until `x` is served, which only a reservation + // on `w0` can do. `big` contributes 10 000 free cpus but has no "foo": `x` never runs there, so + // narrow work fills it unconditionally and more narrow tasks wait than it can take. + let mut rt = TestEnv::new(); + rt.new_named_resource("foo"); + let big = rt.new_worker(&WorkerBuilder::new(10_000)); + let w0 = rt.new_worker(&WorkerBuilder::new(8).res_sum("foo", 1)); + let w1 = rt.new_worker(&WorkerBuilder::new(8).res_sum("foo", 1)); + rt.new_task_running(&TaskBuilder::new().cpus(4), w0); + rt.new_task_running(&TaskBuilder::new().cpus(6), w1); + rt.new_task( + &TaskBuilder::new() + .cpus(8) + .add_resource(1, 1) + .user_priority(10), + ); + let narrow = rt.new_tasks(10_020, &TaskBuilder::new().cpus(1).user_priority(5)); + rt.schedule(); + + assert_eq!( + narrow_per_worker(&rt, &narrow, &[big, w0, w1]), + vec![10_000, 0, 2], + "a reservation on w0 must release w1 even when placements are worth little" + ); +} + +#[test] +fn test_schedule_one_worker_cannot_be_reserved_for_two_blockers() { + // `x` and `y` are different requests at equal priority, both needing the single "foo" of a + // worker, and neither fits anywhere now. `w0` has 7 free cpus but its "foo" is taken; `w1` is + // fully occupied. Neither covers any part of either request, so both are held on the higher + // id, `w1`, which has no free resources. Reservations there consume nothing, so without a + // per-worker limit `x` and `y` are both served by one worker that can later host only one + // of them, and narrow work fills `w0`. With the limit, the one not reserved for is held back + // on `w0` by its own condition, except for the gap: `y` (7 cpus) can never use more than 7 of + // `w0`'s 8 cpus, so one narrow task may still run there. + let mut rt = TestEnv::new(); + rt.new_named_resource("foo"); + let w0 = rt.new_worker(&WorkerBuilder::new(8).res_sum("foo", 1)); + let w1 = rt.new_worker(&WorkerBuilder::new(8).res_sum("foo", 1)); + rt.new_task_running(&TaskBuilder::new().cpus(1).add_resource(1, 1), w0); + rt.new_task_running(&TaskBuilder::new().cpus(8).add_resource(1, 1), w1); + rt.new_task( + &TaskBuilder::new() + .cpus(8) + .add_resource(1, 1) + .user_priority(10), + ); + rt.new_task( + &TaskBuilder::new() + .cpus(7) + .add_resource(1, 1) + .user_priority(10), + ); + let narrow = rt.new_tasks(20, &TaskBuilder::new().cpus(1).user_priority(5)); + rt.schedule(); + + assert_eq!( + narrow_per_worker(&rt, &narrow, &[w0, w1]), + vec![1, 0], + "one worker can be reserved for at most one blocker, so one of them stays unserved" + ); +} + +#[test] +fn test_schedule_reservation_leaves_gap_allowance() { + // A reservation withholds the worker's free resources for the blocker, but part of them is + // capacity the blocker can never use: its gap allowance. The priority condition already lets + // gap tasks run there, and the reservation must not take that capacity away. + // + // Two workers of 16 cpus and 4 gpus each run four 1-cpu + 1-gpu tasks, so every gpu is taken + // and the blocker (2 cpus + 1 gpu) fits nowhere. A full packing of the blocker uses 8 cpus + // and all 4 gpus, so 8 cpus per worker are never usable by it. While k of the running tasks + // remain it fits 4 - k times and cannot use 8 + k cpus, so the gap is 8 cpus (the bound + // `C - M(C) - o` would charge the running tasks' 4 cpus and give only 4). `w1` is held (tie, + // higher id); a reservation there serves the blocker and releases `w0`, but it must consume + // only 12 - 8 = 4 of `w1`'s free cpus, leaving room for 8 gap tasks. + let mut rt = TestEnv::new(); + rt.new_named_resource("gpu"); + let ws = rt.new_workers(2, &WorkerBuilder::new(16).res_sum("gpu", 4)); + for w in &ws { + for _ in 0..4 { + rt.new_task_running(&TaskBuilder::new().cpus(1).add_resource(1, 1), *w); + } + } + rt.new_task( + &TaskBuilder::new() + .cpus(2) + .add_resource(1, 1) + .user_priority(10), + ); + let narrow = rt.new_tasks(40, &TaskBuilder::new().cpus(1).user_priority(5)); + rt.schedule(); + + assert_eq!( + narrow_per_worker(&rt, &narrow, &ws), + vec![12, 8], + "w0 is released by the reservation; the reserved w1 keeps its 8-cpu gap" + ); +} + +#[test] +fn test_schedule_reservation_example() { + // ten 8-cpu workers, each running a single + // 4-cpu task, one 6-cpu task at priority 2 and a hundred 1-cpu tasks at priority 1. The 6-cpu + // task fits nowhere, so one worker is held and reserved for it, which releases the other nine + // (4 cpus each). On the reserved worker the reservation withholds only what the 6-cpu task + // could use: 8 mod 6 = 2 cpus are its gap and stay available. + // + // The single 4-cpu occupant matters: four 1-cpu occupants could leave 6 cpus free at an + // intermediate point, and the gap would be 0. + let mut rt = TestEnv::new(); + let ws = rt.new_workers(10, &WorkerBuilder::new(8)); + for w in &ws { + rt.new_task_running(&TaskBuilder::new().cpus(4), *w); + } + rt.new_task(&TaskBuilder::new().cpus(6).user_priority(2)); + let narrow = rt.new_tasks(100, &TaskBuilder::new().cpus(1).user_priority(1)); + rt.schedule(); + + let mut expected = vec![4; 9]; + expected.push(2); + assert_eq!(narrow_per_worker(&rt, &narrow, &ws), expected); +} + #[test] fn test_schedule_blockers_hold_all_needed_workers() { const N_WORKERS: usize = 4; @@ -2034,3 +2405,161 @@ fn test_reservation_one_assigned() { .sum::(); assert_eq!(filler_scheduled, 6); } + +#[test] +fn test_schedule_gap_tasks_do_not_move_the_hold() { + // The blocker needs 6 cpus. `w0` runs one 4-cpu task: 4 cpus free, 2 of them a gap the blocker + // can never use. `w1` runs one 5-cpu task: 3 cpus free, 2 of them a gap. `w0` is held and its + // gap is filled with 2 narrow tasks. Then 3 narrow tasks finish on `w1`. By plain free cpus + // `w1` (3) now looks better than `w0` (2), but the blocker can use only 1 of them, against 2 on + // `w0`. The hold must stay on `w0`: filling a gap never makes a worker less useful to the + // blocker, so it must not cost the worker its hold (and then its drained capacity). + let mut rt = TestEnv::new(); + let w0 = rt.new_worker(&WorkerBuilder::new(8)); + let w1 = rt.new_worker(&WorkerBuilder::new(8)); + rt.new_task_running(&TaskBuilder::new().cpus(4), w0); + rt.new_task_running(&TaskBuilder::new().cpus(5), w1); + rt.new_task(&TaskBuilder::new().cpus(6).user_priority(10)); + let narrow = rt.new_tasks(40, &TaskBuilder::new().cpus(1).user_priority(5)); + rt.schedule(); + assert_eq!(narrow_per_worker(&rt, &narrow, &[w0, w1]), vec![2, 3]); + + let finished: Vec<_> = narrow + .iter() + .copied() + .filter(|t| { + matches!(rt.task(*t).state, + TaskRuntimeState::Assigned { worker_id, .. } if worker_id == w1) + }) + .collect(); + for t in &finished { + rt.finish_task(*t, w1); + } + let live: Vec<_> = narrow + .iter() + .copied() + .filter(|t| !finished.contains(t)) + .collect(); + rt.schedule(); + assert_eq!( + narrow_per_worker(&rt, &live, &[w0, w1]), + vec![2, 3], + "w0 keeps its hold; w1 is refilled" + ); +} + +#[test] +fn test_schedule_blocker_count_allowance_is_not_multiplied_per_worker() { + let mut rt = TestEnv::new(); + let ws = rt.new_workers(3, &WorkerBuilder::new(12)); + for w in &ws { + rt.new_task_running(&TaskBuilder::new().cpus(5), *w); + } + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(2)); + let above = rt.new_tasks(3, &TaskBuilder::new().cpus(1).user_priority(3)); + let below = rt.new_tasks(40, &TaskBuilder::new().cpus(1).user_priority(1)); + + rt.schedule(); + + let per_worker: Vec = narrow_per_worker(&rt, &above, &ws) + .iter() + .zip(narrow_per_worker(&rt, &below, &ws)) + .map(|(a, b)| a + b) + .collect(); + assert!( + per_worker.iter().any(|placed| *placed <= 4), + "every worker took more than its 4-cpu gap ({per_worker:?}), \ + so the blocker cannot start anywhere once the 5-cpu tasks finish" + ); +} + +#[test] +fn test_schedule_two_requests_share_one_gap_allowance() { + let mut rt = TestEnv::new(); + let foo = rt.new_named_resource("foo"); + let w = rt.new_worker(&WorkerBuilder::new(12).res_sum("foo", 10)); + rt.new_task_running(&TaskBuilder::new().cpus(5), w); + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + let plain = rt.new_tasks(20, &TaskBuilder::new().cpus(1).user_priority(5)); + let with_foo = rt.new_tasks( + 20, + &TaskBuilder::new() + .cpus(1) + .add_resource(foo, 1) + .user_priority(5), + ); + + rt.schedule(); + + let plain_placed = narrow_per_worker(&rt, &plain, &[w])[0]; + let foo_placed = narrow_per_worker(&rt, &with_foo, &[w])[0]; + assert!( + plain_placed + foo_placed <= 4, + "worker runs {} lower-priority cpus ({plain_placed} plain + {foo_placed} with foo), \ + but the blocker leaves a gap of 4 there", + plain_placed + foo_placed + ); +} + +#[test] +fn test_schedule_soft_rejected_worker_is_still_limited_by_the_blocker() { + let mut rt = TestEnv::new(); + let w = rt.new_worker(&WorkerBuilder::new(12)); + rt.new_task_running(&TaskBuilder::new().cpus(5), w); + let blocker = rt.new_task(&TaskBuilder::new().cpus(8).user_priority(10)); + let narrow = rt.new_tasks(20, &TaskBuilder::new().cpus(1).user_priority(1)); + + let blocker_rq = rt.task(blocker).resource_rq_id; + rt.core() + .get_worker_mut(w) + .block_request(blocker_rq, ResourceVariantId::new(0)); + + rt.schedule(); + + assert_eq!( + narrow_per_worker(&rt, &narrow, &[w]), + vec![4], + "a soft-rejected worker keeps its gap limit" + ); +} + +/// The excess of a request over a worker's gap must never exceed what the request actually runs +/// there. The excess lowers the shared gap constraint, so excess that no task uses would buy other +/// requests room beyond the gap, paid from an allowance spent on nothing. +/// +/// `w0` has 12 cpus with a 5-cpu task running: 7 free, and a gap of 4 for the 8-cpu blocker at +/// priority 5. The donor request (2-cpu tasks) has three tasks above the blocker, so it carries an +/// allowance of 3, and they are placed on `w1`, which is too small for the blocker. The donor +/// therefore runs nothing on `w0` while holding an allowance there. Its own tasks may use that +/// allowance and outrank the blocker, but the other two requests have no allowance at all: they +/// share `w0`'s gap and must stay within it together, 4 cpus rather than 4 each. +#[test] +fn test_schedule_unused_allowance_does_not_widen_the_shared_gap() { + let mut rt = TestEnv::new(); + let foo = rt.new_named_resource("foo"); + let w0 = rt.new_worker(&WorkerBuilder::new(12).res_sum("foo", 10)); + let w1 = rt.new_worker(&WorkerBuilder::new(6)); + rt.new_task_running(&TaskBuilder::new().cpus(5), w0); + rt.new_tasks(3, &TaskBuilder::new().cpus(2).user_priority(9)); + rt.new_tasks(10, &TaskBuilder::new().cpus(2).user_priority(1)); + rt.new_task(&TaskBuilder::new().cpus(8).user_priority(5)); + let b1 = rt.new_tasks(20, &TaskBuilder::new().cpus(1).user_priority(1)); + let b2 = rt.new_tasks( + 20, + &TaskBuilder::new() + .cpus(1) + .add_resource(foo, 1) + .user_priority(1), + ); + + rt.schedule(); + + let b1_cpus = narrow_per_worker(&rt, &b1, &[w0])[0]; + let b2_cpus = narrow_per_worker(&rt, &b2, &[w0])[0]; + assert!( + b1_cpus + b2_cpus <= 4, + "requests without an allowance use {} cpus of w0 ({b1_cpus} + {b2_cpus}), \ + but they share a gap of 4", + b1_cpus + b2_cpus + ); +}