Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 16 additions & 7 deletions crates/onnx-runtime-ep-cpu/src/decode_spmd.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,8 +68,10 @@
//! under greedy decode.
//!
//! The worker count is
//! [`crate::kernels::matmul_nbits::configured_persistent_decode_threads`] (about
//! half the logical CPUs); a `THREADS=0` opt-out leaves the decode path unchanged.
//! [`crate::kernels::matmul_nbits::configured_persistent_decode_threads`] (one
//! worker per *allowed physical core*, falling back to half the logical CPUs
//! only when the core topology is undiscoverable); a `THREADS=0` opt-out leaves
//! the decode path unchanged.
//!
//! # Precedence when on (default or `=1`) vs the affinity control
//!
Expand Down Expand Up @@ -1726,8 +1728,9 @@ pub fn shutdown_pools() {

/// Resolve the persistent pool's worker count. Honors `ONNX_GENAI_CPU_DECODE_THREADS`
/// when set (`0` opts out); when unset it uses the persistent-specific default
/// (about half the logical CPUs), *not* the flat pool's eight-worker ceiling --
/// see [`crate::kernels::matmul_nbits::configured_persistent_decode_threads`].
/// (one worker per allowed physical core), *not* the flat pool's eight-worker
/// ceiling -- see
/// [`crate::kernels::matmul_nbits::configured_persistent_decode_threads`].
fn default_threads() -> Option<usize> {
crate::kernels::matmul_nbits::configured_persistent_decode_threads()
}
Expand Down Expand Up @@ -1983,10 +1986,16 @@ fn node_shards(total: usize) -> Vec<NodeShard> {
/// run; reserving *per node* (see [`reserve_split_headroom`]) rather than once
/// globally guarantees the spare core lands on whichever socket the scheduler
/// places the dispatcher on, and keeps every node's completion-counter reads
/// unblocked on a NUMA-split layout. One is enough because the workers already
/// plateau around half the logical CPUs (memory-bandwidth bound), so giving up
/// the single highest-index core costs nothing measurable while removing the
/// unblocked on a NUMA-split layout. One is enough because there is exactly one
/// inline dispatcher thread to house, so giving up the single highest-index core
/// costs one worker's share of a bandwidth-bound loop while removing the
/// starvation cliff.
///
/// This used to be justified by the workers plateauing "around half the logical
/// CPUs", which stopped being the sizing rule when the default became one worker
/// per allowed physical core: on a non-SMT host the pool is now deliberately
/// wide enough that this reserve is what keeps it from being *fully*
/// subscribed.
const DISPATCHER_RESERVED_CPUS: usize = 1;

/// Cap a single pinned worker group so at least [`DISPATCHER_RESERVED_CPUS`]
Expand Down
126 changes: 101 additions & 25 deletions crates/onnx-runtime-ep-cpu/src/kernels/matmul_nbits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3706,8 +3706,8 @@ fn configured_decode_threads() -> Option<usize> {
///
/// It honors `ONNX_GENAI_CPU_DECODE_THREADS` when set (`0` opts out), but when
/// the variable is unset it uses a *different, higher* default than the flat
/// pool: [`default_persistent_threads`] (about half the logical CPUs) instead of
/// the flat pool's eight-worker ceiling. The flat Rayon pool caps at eight
/// pool: [`default_persistent_threads`] (one worker per allowed physical core)
/// instead of the flat pool's eight-worker ceiling. The flat Rayon pool caps at eight
/// because its per-op fork/join regresses beyond that; the persistent pool
/// replaces that fork/join with one hot broadcast barrier, so it keeps scaling
/// with cores until it hits the memory-bandwidth knee (measured plateau ~half
Expand All @@ -3721,6 +3721,7 @@ pub fn configured_persistent_decode_threads() -> Option<usize> {
decode_threads_override(),
value.as_deref(),
available,
crate::core_topology::allowed_physical_cores(),
);
// Snapdragon/X Elite style ARM64 hosts have measured their decode roofline
// at 6--8 workers, and the KAI-style packed SDOT path still scales from the
Expand Down Expand Up @@ -4399,10 +4400,54 @@ fn prefill_worker_name(index: usize) -> String {
/// *spin* before parking, a fully-subscribed pool starves the dispatcher and
/// collapses throughput (measured 1.4 tok/s at 96 workers vs 28.7 at 48 on a
/// 96-logical-CPU host); half sits at the measured plateau while avoiding that
/// cliff, and on SMT hosts it maps to roughly the physical-core count.
fn default_persistent_threads(available: usize) -> Option<usize> {
/// cliff.
///
/// # Why `available / 2` is only a proxy
///
/// The rule the measurements actually support is *one worker per physical
/// core*. Halving the CPU count is a stand-in that agrees with it only when the
/// process can see every logical CPU on an SMT host. It disagrees badly the
/// moment the allowed set is already one CPU per core: `taskset -c 0,2,...,30`
/// on a 16-core/32-thread host leaves `available == 16`, all of them distinct
/// cores, and halving builds an **8**-worker pool on 16 reserved cores. The
/// operator asked for the machine and got half of it, silently, because the
/// proxy cannot tell a 16-CPU SMT-free set from half of a 32-CPU SMT set.
///
/// So take the real answer when the topology is discoverable and the allowed
/// set is known ([`crate::core_topology::allowed_physical_cores`]), and fall
/// back to the halving proxy only when it is not. `None` from the topology
/// means "unknown", never "one" — capping on a guess is how a 32-thread host
/// silently becomes a 1-thread host.
///
/// # What this is *not* evidenced by
///
/// Both measurements behind the old rule — the plateau at "~half the logical
/// CPUs" on a 2-socket Xeon 8480C, and the 1.4 tok/s at 96 workers against 28.7
/// at 48 on a 96-logical-CPU host — were taken on SMT hosts, where half the
/// logical CPUs *is* the physical-core count. Neither one distinguishes the two
/// rules, so neither one supports this change on a **non-SMT** host, where
/// `cores == available` and the default moves from half the machine to all of
/// it. That case is an inference, not a measurement, and no non-SMT host was
/// available to check it.
///
/// What makes the inference safe rather than merely plausible is that the
/// consuming path reserves dispatcher headroom:
/// `decode_spmd`'s `reserve_single_group_headroom` turns a request for `n`
/// on `n` cores into `n - 1` spawned workers plus the inline dispatcher's own
/// shard, so the fully-subscribed spinning layout that produced the 1.4 tok/s
/// cliff is not reachable through this default. Verified on a 16-core host
/// pinned to one CPU per core: 15 threads named `onnx-genai-spmd` on 15 leader
/// CPUs, with the sixteenth core left free for the dispatcher.
fn default_persistent_threads(
available: usize,
allowed_physical_cores: Option<usize>,
) -> Option<usize> {
let available = std::num::NonZeroUsize::new(available)?.get();
Some((available / 2).max(1))
let default = match allowed_physical_cores {
Some(cores) if cores > 0 => cores,
_ => available / 2,
};
Some(default.clamp(1, available))
}

/// Query the performance-core count on Apple Silicon via sysctl.
Expand Down Expand Up @@ -4436,16 +4481,17 @@ fn performance_core_count() -> Option<usize> {
/// unparseable value falls back to [`default_persistent_threads`].
#[cfg(test)]
fn resolve_persistent_decode_threads(raw: Option<&str>, available: usize) -> Option<usize> {
resolve_persistent_decode_threads_with_override(None, raw, available)
resolve_persistent_decode_threads_with_override(None, raw, available, None)
}

fn resolve_persistent_decode_threads_with_override(
override_threads: Option<usize>,
raw: Option<&str>,
available: usize,
allowed_physical_cores: Option<usize>,
) -> Option<usize> {
let available = std::num::NonZeroUsize::new(available)?.get();
let default = default_persistent_threads(available)?;
let default = default_persistent_threads(available, allowed_physical_cores)?;
let threads = match override_threads {
Some(threads) => threads,
None => match raw {
Expand Down Expand Up @@ -4782,18 +4828,22 @@ thread_local! {
/// The lazily built `numa-split` decode layout, or `None` when the mode is not
/// requested or the host cannot be split (fallback, logged once).
///
/// It is sized from [`configured_persistent_decode_threads`] (about half the
/// logical CPUs), *not* the flat pool's eight-worker ceiling. `numa-split` is
/// It is sized from [`configured_persistent_decode_threads`] (one worker per
/// allowed physical core), *not* the flat pool's eight-worker ceiling. `numa-split` is
/// the two-level, node-pinned mirror of the persistent SPMD pool (see
/// [`crate::decode_spmd`]) and its whole purpose is to reach *both* sockets'
/// memory bandwidth; the eight-worker flat ceiling would leave only ~four
/// row-sharded workers per node, far too few to saturate either memory
/// controller, so it could never realize the bandwidth win the layout exists
/// for. Half the logical CPUs, split across the nodes, lands each per-node
/// sub-pool at the measured bandwidth knee while leaving cores for the
/// dispatcher and co-tenants (a *fully*-subscribed split oversubscribes the
/// cores and collapses throughput). `ONNX_GENAI_CPU_DECODE_THREADS` still
/// overrides the count (and `0` opts out).
/// for. One worker per physical core, split across the nodes, lands each
/// per-node sub-pool at the measured bandwidth knee; [`reserve_split_headroom`]
/// is what then leaves room for the dispatcher, since a *fully*-subscribed split
/// oversubscribes the cores and collapses throughput.
/// `ONNX_GENAI_CPU_DECODE_THREADS` still overrides the count (and `0` opts out).
///
/// Note the per-node reserve frees one *logical* CPU rather than a physical
/// core, so on an SMT host the on-node dispatcher gets a sibling rather than a
/// free core; tracked separately as issue 1791.
fn numa_pools() -> Option<&'static crate::decode_numa::NumaDecodePools> {
static NUMA_POOLS: OnceLock<Option<crate::decode_numa::NumaDecodePools>> = OnceLock::new();
NUMA_POOLS
Expand Down Expand Up @@ -16917,18 +16967,44 @@ mod tests {
}

#[test]
fn persistent_decode_thread_default_is_half_the_logical_cpus() {
fn persistent_decode_thread_default_is_half_the_logical_cpus_when_topology_is_unknown() {
// The persistent pool scales past the flat pool's 8-worker ceiling: unset
// -> half the logical CPUs (topology-derived, rule 2), not the flat cap.
assert_eq!(default_persistent_threads(96), Some(48));
assert_eq!(default_persistent_threads(8), Some(4));
assert_eq!(default_persistent_threads(4), Some(2));
assert_eq!(default_persistent_threads(2), Some(1));
assert_eq!(default_persistent_threads(1), Some(1));
assert_eq!(default_persistent_threads(0), None);
// With no discoverable core count the halving proxy still applies.
assert_eq!(default_persistent_threads(96, None), Some(48));
assert_eq!(default_persistent_threads(8, None), Some(4));
assert_eq!(default_persistent_threads(4, None), Some(2));
assert_eq!(default_persistent_threads(2, None), Some(1));
assert_eq!(default_persistent_threads(1, None), Some(1));
assert_eq!(default_persistent_threads(0, None), None);
// Distinct from the flat default on a big host (48 vs 8) -- proving the
// persistent path does not inherit the fork/join-bound cap.
assert_ne!(default_persistent_threads(96), default_decode_threads(96));
assert_ne!(
default_persistent_threads(96, None),
default_decode_threads(96)
);
}

/// The halving proxy is wrong whenever the allowed set is *already* one CPU
/// per physical core, which is what `taskset -c 0,2,...` produces on an SMT
/// host and what every careful benchmark on such a host uses.
#[test]
fn persistent_decode_default_uses_physical_cores_not_half_the_cpuset() {
// 16 CPUs that are 16 distinct cores: the whole point of the pin is to
// get 16 workers. Halving would build 8 and report nothing.
assert_eq!(default_persistent_threads(16, Some(16)), Some(16));
// The unpinned SMT case must be unchanged: 32 CPUs, 16 cores -> 16,
// which is exactly what the halving proxy used to produce.
assert_eq!(default_persistent_threads(32, Some(16)), Some(16));
assert_eq!(default_persistent_threads(96, Some(48)), Some(48));
// A core count can never exceed the CPUs actually available, and can
// never round down to zero.
assert_eq!(default_persistent_threads(4, Some(9)), Some(4));
// `Some(0)` is nonsense from the topology layer; treat it as unknown
// rather than building a zero-worker pool.
assert_eq!(default_persistent_threads(8, Some(0)), Some(4));
// Unknown topology keeps the historical behaviour exactly.
assert_eq!(default_persistent_threads(16, None), Some(8));
}

#[test]
Expand Down Expand Up @@ -17012,15 +17088,15 @@ mod tests {
Some(6)
);
assert_eq!(
resolve_persistent_decode_threads_with_override(Some(8), Some("32"), 96),
resolve_persistent_decode_threads_with_override(Some(8), Some("32"), 96, None),
Some(8)
);
assert_eq!(
resolve_dense_decode_threads_with_override(Some(12), Some("0"), 96),
Some(12)
);
assert_eq!(
resolve_persistent_decode_threads_with_override(None, None, 96),
resolve_persistent_decode_threads_with_override(None, None, 96, None),
Some(48),
"the uncapped automatic default must remain unchanged"
);
Expand All @@ -17045,7 +17121,7 @@ mod tests {
assert_ne!(default_dense_decode_threads(96), default_decode_threads(96));
assert_ne!(
default_dense_decode_threads(96),
default_persistent_threads(96)
default_persistent_threads(96, None)
);
}

Expand Down
Loading