This repository has no description
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}