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

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.

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.

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.

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