This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-resource / src / cpu.rs
6.4 kB 258 lines
1use std::sync::OnceLock; 2use std::sync::atomic::{AtomicUsize, Ordering}; 3 4#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] 5pub struct ThreadCount(usize); 6 7impl ThreadCount { 8 pub const fn new(count: usize) -> Self { 9 Self(if count == 0 { 1 } else { count }) 10 } 11 12 pub const fn get(self) -> usize { 13 self.0 14 } 15} 16 17knot_types::scalar_newtype! { 18 struct WorkUnits(usize); 19} 20 21struct Budget { 22 ceiling: usize, 23 available: AtomicUsize, 24} 25 26static BUDGET: OnceLock<Budget> = OnceLock::new(); 27 28fn detected_threads() -> ThreadCount { 29 ThreadCount::new( 30 std::thread::available_parallelism() 31 .map(std::num::NonZeroUsize::get) 32 .unwrap_or(1), 33 ) 34} 35 36pub(crate) fn install(ceiling: ThreadCount) -> ThreadCount { 37 let budget = BUDGET.get_or_init(|| Budget { 38 ceiling: ceiling.get(), 39 available: AtomicUsize::new(ceiling.get().saturating_sub(1)), 40 }); 41 ThreadCount::new(budget.ceiling) 42} 43 44fn budget() -> &'static Budget { 45 BUDGET.get_or_init(|| { 46 let ceiling = detected_threads().get(); 47 Budget { 48 ceiling, 49 available: AtomicUsize::new(ceiling.saturating_sub(1)), 50 } 51 }) 52} 53 54pub fn threads() -> ThreadCount { 55 ThreadCount::new(budget().ceiling) 56} 57 58pub fn gix_thread_limit() -> Option<ThreadCount> { 59 let budget = budget(); 60 (budget.ceiling < detected_threads().get()).then_some(ThreadCount::new(budget.ceiling)) 61} 62 63pub(crate) fn ceiling() -> usize { 64 budget().ceiling 65} 66 67static SATURATE: AtomicUsize = AtomicUsize::new(0); 68 69pub struct Saturate(()); 70 71impl Drop for Saturate { 72 fn drop(&mut self) { 73 SATURATE.fetch_sub(1, Ordering::Relaxed); 74 } 75} 76 77// the eat my machine button 78pub fn saturate() -> Saturate { 79 SATURATE.fetch_add(1, Ordering::Relaxed); 80 Saturate(()) 81} 82 83struct Lease { 84 extra: usize, 85 saturated: bool, 86} 87 88impl Lease { 89 fn none() -> Self { 90 Self { 91 extra: 0, 92 saturated: false, 93 } 94 } 95 96 fn lanes(&self) -> usize { 97 self.extra + 1 98 } 99} 100 101impl Drop for Lease { 102 fn drop(&mut self) { 103 if !self.saturated && self.extra > 0 { 104 budget().available.fetch_add(self.extra, Ordering::AcqRel); 105 } 106 } 107} 108 109fn lease(units: WorkUnits) -> Lease { 110 let budget = budget(); 111 let want = units 112 .get() 113 .saturating_sub(1) 114 .min(budget.ceiling.saturating_sub(1)); 115 if want == 0 { 116 return Lease::none(); 117 } 118 if SATURATE.load(Ordering::Relaxed) > 0 { 119 return Lease { 120 extra: want, 121 saturated: true, 122 }; 123 } 124 let mut available = budget.available.load(Ordering::Relaxed); 125 loop { 126 let grant = want.min(available); 127 if grant == 0 { 128 return Lease::none(); 129 } 130 match budget.available.compare_exchange_weak( 131 available, 132 available - grant, 133 Ordering::AcqRel, 134 Ordering::Relaxed, 135 ) { 136 Ok(_) => { 137 return Lease { 138 extra: grant, 139 saturated: false, 140 }; 141 } 142 Err(observed) => available = observed, 143 } 144 } 145} 146 147pub fn map_spans<R, E, F>(len: usize, f: F) -> Result<Vec<R>, E> 148where 149 R: Send, 150 E: Send, 151 F: Fn(usize, usize) -> Result<Vec<R>, E> + Sync, 152{ 153 if len == 0 { 154 return Ok(Vec::new()); 155 } 156 let f = &f; 157 let lease = lease(WorkUnits::new(len)); 158 let lanes = lease.lanes().min(len); 159 if lanes <= 1 { 160 return f(0, len); 161 } 162 let chunk = len.div_ceil(lanes); 163 let spans: Vec<(usize, usize)> = (0..lanes) 164 .map(|lane| (lane * chunk, ((lane + 1) * chunk).min(len))) 165 .filter(|(start, end)| start < end) 166 .collect(); 167 let (head, tail) = spans 168 .split_first() 169 .expect("a positive lane count yields at least one span"); 170 let ordered: Vec<Result<Vec<R>, E>> = std::thread::scope(|scope| { 171 let handles: Vec<_> = tail 172 .iter() 173 .map(|&(start, end)| scope.spawn(move || f(start, end))) 174 .collect(); 175 let head = f(head.0, head.1); 176 std::iter::once(head) 177 .chain( 178 handles 179 .into_iter() 180 .map(|handle| handle.join().expect("resource worker panicked")), 181 ) 182 .collect() 183 }); 184 ordered.into_iter().try_fold(Vec::new(), |mut acc, part| { 185 acc.extend(part?); 186 Ok(acc) 187 }) 188} 189 190pub fn map_chunks<T, R, E, F>(items: &[T], f: F) -> Result<Vec<R>, E> 191where 192 T: Sync, 193 R: Send, 194 E: Send, 195 F: Fn(&[T]) -> Result<Vec<R>, E> + Sync, 196{ 197 map_spans(items.len(), |start, end| f(&items[start..end])) 198} 199 200#[cfg(test)] 201mod tests { 202 use super::*; 203 204 #[test] 205 fn a_positive_span_count_covers_every_index_in_order() { 206 let doubled = map_spans::<usize, (), _>(1000, |start, end| { 207 Ok((start..end).map(|value| value * 2).collect()) 208 }) 209 .unwrap(); 210 assert_eq!(doubled.len(), 1000); 211 assert!( 212 doubled 213 .iter() 214 .enumerate() 215 .all(|(index, value)| *value == index * 2) 216 ); 217 } 218 219 #[test] 220 fn an_empty_span_produces_nothing() { 221 assert_eq!( 222 map_spans::<usize, (), _>(0, |_, _| Ok(vec![1])).unwrap(), 223 Vec::<usize>::new() 224 ); 225 } 226 227 #[test] 228 fn a_worker_error_propagates() { 229 let outcome = 230 map_spans::<usize, &str, _>( 231 500, 232 |start, _| { 233 if start == 0 { Err("boom") } else { Ok(vec![]) } 234 }, 235 ); 236 assert_eq!(outcome, Err("boom")); 237 } 238 239 #[test] 240 fn chunks_preserve_element_order() { 241 let items: Vec<usize> = (0..777).collect(); 242 let echoed = map_chunks::<usize, usize, (), _>(&items, |batch| Ok(batch.to_vec())).unwrap(); 243 assert_eq!(echoed, items); 244 } 245 246 #[test] 247 fn a_lease_never_reserves_more_than_the_ceiling() { 248 let lease = lease(WorkUnits::new(usize::MAX)); 249 assert!(lease.lanes() <= threads().get()); 250 } 251 252 #[test] 253 fn a_saturated_lease_still_respects_the_ceiling() { 254 let _boost = saturate(); 255 let lease = lease(WorkUnits::new(usize::MAX)); 256 assert!(lease.lanes() <= threads().get()); 257 } 258}