pub struct Pool { /* private fields */ }Expand description
Persistent thread pool: shared job slot, epoch dispatch, caller participation.
Implementations§
Source§impl Pool
impl Pool
Sourcepub fn new(n_workers: usize) -> Self
pub fn new(n_workers: usize) -> Self
Examples found in repository?
13fn main() {
14 for nt in [2usize, 4, 6, 8, 10] {
15 let pool = Pool::new(nt);
16 let noop: &(dyn Fn(usize, usize) + Sync) = &|_w, _n| {};
17 // Warm: first dispatch spawns/parks the workers.
18 for _ in 0..100 {
19 pool.run(noop);
20 }
21 let iters = 20_000;
22 let t0 = Instant::now();
23 for _ in 0..iters {
24 pool.run(noop);
25 }
26 let el = t0.elapsed().as_secs_f64();
27 let per = el / iters as f64 * 1e6;
28 println!(
29 "threads={nt:2} {per:7.2} us/dispatch -> {:6.2} ms/token at 200 matvecs",
30 per * 200.0 / 1e3
31 );
32 }
33}More examples
21fn main() {
22 let mut args = std::env::args().skip(1);
23 let path = args
24 .next()
25 .expect("usage: matvec_bw <model.cmf> [sweep|one <tensor>]");
26 let mode = args.next().unwrap_or_else(|| "sweep".to_string());
27 let model = Arc::new(cortiq_core::CmfModel::open(&path).expect("open model"));
28
29 // Real LM activations carry a few heavy channels (>8·rms); measured
30 // mean on this model is ~3.7. NOUT models that distribution.
31 let nout: usize = std::env::var("NOUT")
32 .ok()
33 .and_then(|v| v.parse().ok())
34 .unwrap_or(4);
35 let mk_x = |cols: usize| -> Vec<f32> {
36 let mut x: Vec<f32> = (0..cols).map(|i| ((i % 17) as f32 - 8.0) / 8.0).collect();
37 for k in 0..nout {
38 x[k * 37 % cols] = 40.0;
39 }
40 x
41 };
42
43 if mode == "one" {
44 let name = args
45 .next()
46 .unwrap_or_else(|| "model.embed_tokens.weight".to_string());
47 let entry = model.tensor(&name).expect("tensor not found");
48 let (rows, cols) = (entry.shape[0], entry.shape[1]);
49 let nbytes = entry.nbytes as f64;
50 println!(
51 "tensor {name}: {rows}x{cols} {:?} = {:.1} MB, NOUT={nout}",
52 entry.dtype,
53 nbytes / 1e6
54 );
55 let t = QTensor::from_model(&model, &name).expect("wrap");
56 let x = mk_x(cols);
57 let mut out = vec![0f32; rows];
58 t.matvec(&x, &mut out, None);
59 for nt in [1usize, 2, 4, 6, 8, 10] {
60 let pool = if nt == 1 { None } else { Some(Pool::new(nt)) };
61 let iters = 8;
62 t.matvec(&x, &mut out, pool.as_ref());
63 let t0 = Instant::now();
64 for _ in 0..iters {
65 t.matvec(&x, &mut out, pool.as_ref());
66 }
67 let el = t0.elapsed().as_secs_f64();
68 println!(
69 "threads={nt:2} {:6.2} ms/matvec {:6.1} GB/s (sink {:.3})",
70 el / iters as f64 * 1e3,
71 nbytes * iters as f64 / el / 1e9,
72 out[0]
73 );
74 }
75 return;
76 }
77
78 // sweep: every 2-D q8_2f tensor once = one decode's worth of weights.
79 let names: Vec<String> = model
80 .tensors
81 .iter()
82 .filter(|t| t.dtype == TensorDtype::Q8_2f && t.shape.len() == 2)
83 .map(|t| t.name.clone())
84 .collect();
85 let total_bytes: f64 = model
86 .tensors
87 .iter()
88 .filter(|t| t.dtype == TensorDtype::Q8_2f && t.shape.len() == 2)
89 .map(|t| t.nbytes as f64)
90 .sum();
91 println!(
92 "sweep: {} q8_2f tensors, {:.2} GB total (= weights streamed per decode token), NOUT={nout}",
93 names.len(),
94 total_bytes / 1e9
95 );
96
97 let tensors: Vec<(QTensor, Vec<f32>, Vec<f32>)> = names
98 .iter()
99 .map(|n| {
100 let e = model.tensor(n).unwrap();
101 let (rows, cols) = (e.shape[0], e.shape[1]);
102 (
103 QTensor::from_model(&model, n).expect("wrap"),
104 mk_x(cols),
105 vec![0f32; rows],
106 )
107 })
108 .collect();
109 let mut tensors = tensors;
110
111 // Whole-model residency pass BEFORE any timing: the first touch of a
112 // 4.2 GB mmap faults ~260k pages, which would otherwise be charged to
113 // whichever thread count happens to run first. REVERSE=1 flips the
114 // order as a check that no first-touch cost is left in the table.
115 for _ in 0..2 {
116 for (t, x, out) in tensors.iter_mut() {
117 t.matvec(x, out, None);
118 }
119 }
120 let mut counts = vec![1usize, 2, 4, 6, 8, 10];
121 if std::env::var("REVERSE").is_ok() {
122 counts.reverse();
123 }
124 for nt in counts {
125 let pool = if nt == 1 { None } else { Some(Pool::new(nt)) };
126 // one warm pass, then two measured
127 for (t, x, out) in tensors.iter_mut() {
128 t.matvec(x, out, pool.as_ref());
129 }
130 let iters = 2;
131 let t0 = Instant::now();
132 for _ in 0..iters {
133 for (t, x, out) in tensors.iter_mut() {
134 t.matvec(x, out, pool.as_ref());
135 }
136 }
137 let el = t0.elapsed().as_secs_f64();
138 let per_tok = el / iters as f64;
139 println!(
140 "threads={nt:2} {:7.1} ms/sweep {:6.1} GB/s -> weight-path-only ceiling {:5.1} tok/s",
141 per_tok * 1e3,
142 total_bytes * iters as f64 / el / 1e9,
143 1.0 / per_tok
144 );
145 }
146}Sourcepub fn with_spin(n_workers: usize, spin_budget: usize) -> Self
pub fn with_spin(n_workers: usize, spin_budget: usize) -> Self
Explicit spin budget (tests pin it without touching the env).
Sourcepub fn effective_threads() -> usize
pub fn effective_threads() -> usize
The thread count from_env would use RIGHT NOW: forced (C ABI)
CMF_THREADS > big-core topology > available_parallelism−1. ≤1 means the model runs serial (no pool). Introspection (
execution_mode, status endpoints) must report THIS, not available_parallelism.
Examples found in repository?
175fn main() {
176 let args: Vec<String> = std::env::args().collect();
177 let cmd = args.get(1).map(String::as_str).unwrap_or("");
178 let out = PathBuf::from(arg(&args, "--out").expect("--out"));
179 match cmd {
180 "decode" => {
181 let wav = std::fs::read(arg(&args, "--wav").expect("--wav")).unwrap();
182 let w = mimo_audio::decode_wav(&wav).unwrap();
183 let flat: Vec<f32> = w.channels.concat();
184 save_f32(&out, &[w.channels.len(), w.frames()], &flat);
185 println!(
186 "rate {} channels {} frames {}",
187 w.sample_rate,
188 w.channels.len(),
189 w.frames()
190 );
191 }
192 "frontend" => {
193 std::fs::create_dir_all(&out).unwrap();
194 let wav = std::fs::read(arg(&args, "--wav").expect("--wav")).unwrap();
195 let t0 = Instant::now();
196 let w = mimo_audio::decode_wav(&wav).unwrap();
197 save_f32(
198 &out.join("dec.npy"),
199 &[w.channels.len(), w.frames()],
200 &w.channels.concat(),
201 );
202 let chans: Vec<Vec<f32>> = w
203 .channels
204 .iter()
205 .map(|c| mimo_audio::resample_sinc(c, w.sample_rate, mimo_audio::SAMPLE_RATE))
206 .collect();
207 save_f32(
208 &out.join("chan24k.npy"),
209 &[chans.len(), chans[0].len()],
210 &chans.concat(),
211 );
212 let mono = mimo_audio::wav_to_mono_24k(&w).unwrap();
213 save_f32(&out.join("wave24k.npy"), &[mono.len()], &mono);
214 let pool = cortiq_engine::pool::Pool::from_env();
215 let (mel, m) = mimo_audio::log_mel(&mono, pool.as_deref()).unwrap();
216 save_f32(&out.join("mel.npy"), &[m, mimo_audio::N_MELS], &mel);
217 println!(
218 "rate {} channels {} frames {} -> {} samples, {m} mel frames, K {} ({:.3}s)",
219 w.sample_rate,
220 w.channels.len(),
221 w.frames(),
222 mono.len(),
223 mimo_audio::audio_token_count(m, 4),
224 t0.elapsed().as_secs_f64()
225 );
226 }
227 "tower" => {
228 std::fs::create_dir_all(&out).unwrap();
229 let src = PathBuf::from(arg(&args, "--src").expect("--src"));
230 let t0 = Instant::now();
231 let audio = if src.extension().is_some_and(|e| e == "cmf") {
232 let model = Arc::new(cortiq_core::CmfModel::open(&src).expect("open cmf"));
233 MimoAudio::from_model(&model).expect("load towers")
234 } else {
235 MimoAudio::from_hf_dir(&src).expect("load towers")
236 };
237 let t_load = t0.elapsed().as_secs_f64();
238 let (mel, m) = if let Some(mp) = arg(&args, "--mel") {
239 let (shape, v) = load_f32(Path::new(&mp));
240 assert_eq!(shape[1], mimo_audio::N_MELS);
241 (v, shape[0])
242 } else {
243 let wav = std::fs::read(arg(&args, "--wav").expect("--wav or --mel")).unwrap();
244 audio.wav_to_mel(&wav).unwrap()
245 };
246 let t1 = Instant::now();
247 let feats = audio.features(&mel, m).unwrap();
248 let t_feats = t1.elapsed().as_secs_f64();
249 let d = audio.tokenizer.cfg.d_model;
250 let rows = feats.len() / d;
251 save_f32(&out.join("feats.npy"), &[rows, d], &feats);
252 let t2 = Instant::now();
253 let exact = audio.tokenizer.quantize(&feats, rows, false, audio.pool());
254 let t_rvq = t2.elapsed().as_secs_f64();
255 let rounded = audio.tokenizer.quantize(&feats, rows, true, audio.pool());
256 let levels = exact.len() / rows;
257 save_i32(&out.join("codes_exact.npy"), &[rows, levels], &exact);
258 save_i32(&out.join("codes_bf16books.npy"), &[rows, levels], &rounded);
259 let own = mimo_audio::AudioCodes {
260 frames: rows,
261 levels,
262 codes: if audio.bf16_codebooks {
263 rounded.clone()
264 } else {
265 exact.clone()
266 },
267 };
268 let t3 = Instant::now();
269 let emb_own = audio.embed_codes(&own).unwrap();
270 let t_enc = t3.elapsed().as_secs_f64();
271 save_f32(
272 &out.join("embeds_own.npy"),
273 &[emb_own.n_tokens, emb_own.dim],
274 &emb_own.rows,
275 );
276 let fixed = match arg(&args, "--codes") {
277 Some(cp) => {
278 let (shape, v) = load_codes(Path::new(&cp));
279 mimo_audio::AudioCodes {
280 frames: shape[0],
281 levels: shape[1],
282 codes: v,
283 }
284 }
285 None => own.clone(),
286 };
287 let emb = audio.embed_codes(&fixed).unwrap();
288 save_f32(&out.join("embeds.npy"), &[emb.n_tokens, emb.dim], &emb.rows);
289 let k = mimo_audio::audio_token_count(m, audio.encoder.cfg.group);
290 let meta = serde_json::json!({
291 "src": src.display().to_string(),
292 "mel_frames": m,
293 "segments": mimo_audio::segment_lengths(m),
294 "codes": rows,
295 "placeholder_count_K": k,
296 "embed_rows_own": emb_own.n_tokens,
297 "embed_rows_fixed": emb.n_tokens,
298 "load_s": t_load,
299 "features_s": t_feats,
300 "rvq_s": t_rvq,
301 "encoder_s": t_enc,
302 "bf16_codebooks_default": audio.bf16_codebooks,
303 "threads": cortiq_engine::pool::Pool::effective_threads(),
304 });
305 std::fs::write(
306 out.join("tower.json"),
307 serde_json::to_string_pretty(&meta).unwrap(),
308 )
309 .unwrap();
310 println!("{meta}");
311 assert_eq!(emb_own.n_tokens, k, "placeholder count != encoder rows");
312 }
313 "calib" => {
314 let src = PathBuf::from(arg(&args, "--src").expect("--src"));
315 let model = Arc::new(cortiq_core::CmfModel::open(&src).expect("open cmf"));
316 let audio = MimoAudio::from_model(&model).expect("load towers");
317 let dir = PathBuf::from(arg(&args, "--wav-dir").expect("--wav-dir"));
318 let mut wavs: Vec<PathBuf> = std::fs::read_dir(&dir)
319 .unwrap()
320 .filter_map(|e| e.ok().map(|e| e.path()))
321 .filter(|p| p.extension().is_some_and(|e| e == "wav"))
322 .collect();
323 wavs.sort();
324 let t0 = Instant::now();
325 cortiq_engine::gptq_capture::begin(true);
326 let mut frames = 0usize;
327 for w in &wavs {
328 let emb = audio.embed_wav(&std::fs::read(w).unwrap()).unwrap();
329 frames += emb.n_tokens;
330 eprintln!(
331 " {} -> {} rows ({:.0}s)",
332 w.display(),
333 emb.n_tokens,
334 t0.elapsed().as_secs_f64()
335 );
336 }
337 let hess = cortiq_engine::gptq_capture::end();
338 save_hessians(&out, &hess);
339 println!(
340 "{} clips, {frames} LLM rows, {} linears -> {} ({:.0}s)",
341 wavs.len(),
342 hess.len(),
343 out.display(),
344 t0.elapsed().as_secs_f64()
345 );
346 }
347 _ => {
348 eprintln!(
349 "usage: mimo_audio_dump (decode|frontend|tower|calib) --out ... (see the source header)"
350 );
351 std::process::exit(2);
352 }
353 }
354}Sourcepub fn from_env() -> Option<Arc<Self>>
pub fn from_env() -> Option<Arc<Self>>
Pool sized from CMF_THREADS (see module docs). None = serial.
Without the env, heterogeneous ARM defaults to its BIG cores.
Examples found in repository?
175fn main() {
176 let args: Vec<String> = std::env::args().collect();
177 let cmd = args.get(1).map(String::as_str).unwrap_or("");
178 let out = PathBuf::from(arg(&args, "--out").expect("--out"));
179 match cmd {
180 "decode" => {
181 let wav = std::fs::read(arg(&args, "--wav").expect("--wav")).unwrap();
182 let w = mimo_audio::decode_wav(&wav).unwrap();
183 let flat: Vec<f32> = w.channels.concat();
184 save_f32(&out, &[w.channels.len(), w.frames()], &flat);
185 println!(
186 "rate {} channels {} frames {}",
187 w.sample_rate,
188 w.channels.len(),
189 w.frames()
190 );
191 }
192 "frontend" => {
193 std::fs::create_dir_all(&out).unwrap();
194 let wav = std::fs::read(arg(&args, "--wav").expect("--wav")).unwrap();
195 let t0 = Instant::now();
196 let w = mimo_audio::decode_wav(&wav).unwrap();
197 save_f32(
198 &out.join("dec.npy"),
199 &[w.channels.len(), w.frames()],
200 &w.channels.concat(),
201 );
202 let chans: Vec<Vec<f32>> = w
203 .channels
204 .iter()
205 .map(|c| mimo_audio::resample_sinc(c, w.sample_rate, mimo_audio::SAMPLE_RATE))
206 .collect();
207 save_f32(
208 &out.join("chan24k.npy"),
209 &[chans.len(), chans[0].len()],
210 &chans.concat(),
211 );
212 let mono = mimo_audio::wav_to_mono_24k(&w).unwrap();
213 save_f32(&out.join("wave24k.npy"), &[mono.len()], &mono);
214 let pool = cortiq_engine::pool::Pool::from_env();
215 let (mel, m) = mimo_audio::log_mel(&mono, pool.as_deref()).unwrap();
216 save_f32(&out.join("mel.npy"), &[m, mimo_audio::N_MELS], &mel);
217 println!(
218 "rate {} channels {} frames {} -> {} samples, {m} mel frames, K {} ({:.3}s)",
219 w.sample_rate,
220 w.channels.len(),
221 w.frames(),
222 mono.len(),
223 mimo_audio::audio_token_count(m, 4),
224 t0.elapsed().as_secs_f64()
225 );
226 }
227 "tower" => {
228 std::fs::create_dir_all(&out).unwrap();
229 let src = PathBuf::from(arg(&args, "--src").expect("--src"));
230 let t0 = Instant::now();
231 let audio = if src.extension().is_some_and(|e| e == "cmf") {
232 let model = Arc::new(cortiq_core::CmfModel::open(&src).expect("open cmf"));
233 MimoAudio::from_model(&model).expect("load towers")
234 } else {
235 MimoAudio::from_hf_dir(&src).expect("load towers")
236 };
237 let t_load = t0.elapsed().as_secs_f64();
238 let (mel, m) = if let Some(mp) = arg(&args, "--mel") {
239 let (shape, v) = load_f32(Path::new(&mp));
240 assert_eq!(shape[1], mimo_audio::N_MELS);
241 (v, shape[0])
242 } else {
243 let wav = std::fs::read(arg(&args, "--wav").expect("--wav or --mel")).unwrap();
244 audio.wav_to_mel(&wav).unwrap()
245 };
246 let t1 = Instant::now();
247 let feats = audio.features(&mel, m).unwrap();
248 let t_feats = t1.elapsed().as_secs_f64();
249 let d = audio.tokenizer.cfg.d_model;
250 let rows = feats.len() / d;
251 save_f32(&out.join("feats.npy"), &[rows, d], &feats);
252 let t2 = Instant::now();
253 let exact = audio.tokenizer.quantize(&feats, rows, false, audio.pool());
254 let t_rvq = t2.elapsed().as_secs_f64();
255 let rounded = audio.tokenizer.quantize(&feats, rows, true, audio.pool());
256 let levels = exact.len() / rows;
257 save_i32(&out.join("codes_exact.npy"), &[rows, levels], &exact);
258 save_i32(&out.join("codes_bf16books.npy"), &[rows, levels], &rounded);
259 let own = mimo_audio::AudioCodes {
260 frames: rows,
261 levels,
262 codes: if audio.bf16_codebooks {
263 rounded.clone()
264 } else {
265 exact.clone()
266 },
267 };
268 let t3 = Instant::now();
269 let emb_own = audio.embed_codes(&own).unwrap();
270 let t_enc = t3.elapsed().as_secs_f64();
271 save_f32(
272 &out.join("embeds_own.npy"),
273 &[emb_own.n_tokens, emb_own.dim],
274 &emb_own.rows,
275 );
276 let fixed = match arg(&args, "--codes") {
277 Some(cp) => {
278 let (shape, v) = load_codes(Path::new(&cp));
279 mimo_audio::AudioCodes {
280 frames: shape[0],
281 levels: shape[1],
282 codes: v,
283 }
284 }
285 None => own.clone(),
286 };
287 let emb = audio.embed_codes(&fixed).unwrap();
288 save_f32(&out.join("embeds.npy"), &[emb.n_tokens, emb.dim], &emb.rows);
289 let k = mimo_audio::audio_token_count(m, audio.encoder.cfg.group);
290 let meta = serde_json::json!({
291 "src": src.display().to_string(),
292 "mel_frames": m,
293 "segments": mimo_audio::segment_lengths(m),
294 "codes": rows,
295 "placeholder_count_K": k,
296 "embed_rows_own": emb_own.n_tokens,
297 "embed_rows_fixed": emb.n_tokens,
298 "load_s": t_load,
299 "features_s": t_feats,
300 "rvq_s": t_rvq,
301 "encoder_s": t_enc,
302 "bf16_codebooks_default": audio.bf16_codebooks,
303 "threads": cortiq_engine::pool::Pool::effective_threads(),
304 });
305 std::fs::write(
306 out.join("tower.json"),
307 serde_json::to_string_pretty(&meta).unwrap(),
308 )
309 .unwrap();
310 println!("{meta}");
311 assert_eq!(emb_own.n_tokens, k, "placeholder count != encoder rows");
312 }
313 "calib" => {
314 let src = PathBuf::from(arg(&args, "--src").expect("--src"));
315 let model = Arc::new(cortiq_core::CmfModel::open(&src).expect("open cmf"));
316 let audio = MimoAudio::from_model(&model).expect("load towers");
317 let dir = PathBuf::from(arg(&args, "--wav-dir").expect("--wav-dir"));
318 let mut wavs: Vec<PathBuf> = std::fs::read_dir(&dir)
319 .unwrap()
320 .filter_map(|e| e.ok().map(|e| e.path()))
321 .filter(|p| p.extension().is_some_and(|e| e == "wav"))
322 .collect();
323 wavs.sort();
324 let t0 = Instant::now();
325 cortiq_engine::gptq_capture::begin(true);
326 let mut frames = 0usize;
327 for w in &wavs {
328 let emb = audio.embed_wav(&std::fs::read(w).unwrap()).unwrap();
329 frames += emb.n_tokens;
330 eprintln!(
331 " {} -> {} rows ({:.0}s)",
332 w.display(),
333 emb.n_tokens,
334 t0.elapsed().as_secs_f64()
335 );
336 }
337 let hess = cortiq_engine::gptq_capture::end();
338 save_hessians(&out, &hess);
339 println!(
340 "{} clips, {frames} LLM rows, {} linears -> {} ({:.0}s)",
341 wavs.len(),
342 hess.len(),
343 out.display(),
344 t0.elapsed().as_secs_f64()
345 );
346 }
347 _ => {
348 eprintln!(
349 "usage: mimo_audio_dump (decode|frontend|tower|calib) --out ... (see the source header)"
350 );
351 std::process::exit(2);
352 }
353 }
354}Sourcepub fn bind_numa(&self, regions: &[&[u8]])
pub fn bind_numa(&self, regions: &[&[u8]])
Keep the pool on the NUMA node that holds regions (the model’s
weight bytes). Linux with two or more nodes only; CMF_NUMA=0
turns it off, CMF_NUMA=node:<n> forces a node.
WHY: decode streams every weight once per token, and on a two-socket host the page cache holds a file on whichever node read it. Unpinned, the scheduler spreads the workers over both sockets and half the matvec rows cross the socket link. Measured on a 2×EPYC 7763 pod with the model’s pages all on node 0 (31 CPUs of cgroup quota): a STREAM-style read over a node-0 buffer gives 42 GB/s from 31 unpinned threads and 74 GB/s from 31 threads kept on node 0. The mask is the node’s physical cores (first SMT sibling) when there are enough of them for the pool, else the whole node; never narrower than the pool, so nothing oversubscribes. Threads are bound to a SET of cores, not to one core each: the scheduler still balances inside the node. The calling thread adopts the same mask on its next dispatch.
Sourcepub fn run_rows(&self, rows: usize, f: &(dyn Fn(usize, usize) + Sync))
pub fn run_rows(&self, rows: usize, f: &(dyn Fn(usize, usize) + Sync))
Run f(row_start, row_end) over 0..rows, self-balancing.
One dispatch, but workers pull row-ranges from a shared cursor instead of each taking a fixed 1/n slice. On a heterogeneous CPU (Apple Silicon: 4 P-cores + 6 E-cores here) a static split makes every matvec end at the SLOWEST core’s pace while the fast ones idle at the barrier; pulling by grain lets a P-core take several chunks for each one an E-core takes, so skew collapses to a single grain. Row ranges stay disjoint and each row’s dot is computed exactly as in the serial path → bit-identical output.
Sourcepub fn run_many(&self, parts: &[(usize, &(dyn Fn(usize, usize) + Sync))])
pub fn run_many(&self, parts: &[(usize, &(dyn Fn(usize, usize) + Sync))])
Multi-matrix job: one dispatch serves SEVERAL row spaces
(roadmap §3 P0 — «одна внешняя публикация job на слой»). Parts
are laid out back-to-back in a virtual row space and pulled by
grain from one shared cursor, so QKV or gate+up cost a single
barrier instead of one each. Each part’s f(start, end) sees its
OWN row indices — per-row math and outputs are bit-identical to
separate run_rows calls.
Sourcepub fn run(&self, f: &(dyn Fn(usize, usize) + Sync))
pub fn run(&self, f: &(dyn Fn(usize, usize) + Sync))
Run f(worker_idx, n_participants) on every worker AND the
calling thread (worker_idx = n_workers() for the caller);
returns when all participants have finished.
Examples found in repository?
13fn main() {
14 for nt in [2usize, 4, 6, 8, 10] {
15 let pool = Pool::new(nt);
16 let noop: &(dyn Fn(usize, usize) + Sync) = &|_w, _n| {};
17 // Warm: first dispatch spawns/parks the workers.
18 for _ in 0..100 {
19 pool.run(noop);
20 }
21 let iters = 20_000;
22 let t0 = Instant::now();
23 for _ in 0..iters {
24 pool.run(noop);
25 }
26 let el = t0.elapsed().as_secs_f64();
27 let per = el / iters as f64 * 1e6;
28 println!(
29 "threads={nt:2} {per:7.2} us/dispatch -> {:6.2} ms/token at 200 matvecs",
30 per * 200.0 / 1e3
31 );
32 }
33}