Skip to main content

Pool

Struct Pool 

Source
pub struct Pool { /* private fields */ }
Expand description

Persistent thread pool: shared job slot, epoch dispatch, caller participation.

Implementations§

Source§

impl Pool

Source

pub fn new(n_workers: usize) -> Self

Examples found in repository?
examples/pool_lat.rs (line 15)
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
Hide additional examples
examples/matvec_bw.rs (line 60)
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}
Source

pub fn with_spin(n_workers: usize, spin_budget: usize) -> Self

Explicit spin budget (tests pin it without touching the env).

Source

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?
examples/mimo_audio_dump.rs (line 303)
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}
Source

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?
examples/mimo_audio_dump.rs (line 214)
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}
Source

pub fn n_workers(&self) -> usize

Spawned worker threads (the caller joins each job on top).

Source

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.

Source

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.

Source

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.

Source

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?
examples/pool_lat.rs (line 19)
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}

Trait Implementations§

Source§

impl Drop for Pool

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

§

impl !RefUnwindSafe for Pool

§

impl !UnwindSafe for Pool

§

impl Freeze for Pool

§

impl Send for Pool

§

impl Sync for Pool

§

impl Unpin for Pool

§

impl UnsafeUnpin for Pool

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more