diff --git a/crates/onnx-runtime-ep-cpu/src/decode_affinity.rs b/crates/onnx-runtime-ep-cpu/src/decode_affinity.rs index 9323e6ad69..952577e80e 100644 --- a/crates/onnx-runtime-ep-cpu/src/decode_affinity.rs +++ b/crates/onnx-runtime-ep-cpu/src/decode_affinity.rs @@ -37,6 +37,7 @@ //! unpinned (logged once). use std::collections::BTreeMap; +use std::sync::OnceLock; /// Selects how the decode pool binds its workers to CPUs. /// @@ -55,6 +56,100 @@ pub const DECODE_AFFINITY_ENV: &str = "ONNX_GENAI_CPU_DECODE_AFFINITY"; /// rejected value always sees the full menu of valid options. const ACCEPTED_MODES: &str = "`off`, `compact`, `node:`, `numa-split`"; +/// Selects how decode workers are laid out across the CPUs they are allowed. +/// +/// Orthogonal to [`DECODE_AFFINITY_ENV`], which selects *which* CPUs the pool +/// may use. This selects how the workers are arranged inside that set: +/// +/// * `spread` -- worker `i` takes a distinct physical core for as long as cores +/// last, only then doubling up on SMT siblings. Fastest on a host the process +/// has to itself, which is what `crate::core_topology`'s module docs measure. +/// * `compact` -- workers fill a core's SMT siblings before moving to the next +/// core, so an N-worker pool occupies about `N / siblings` cores and leaves +/// the rest of the machine free for a co-tenant. +/// +/// Neither is universally better, which is exactly why it is a selector and not +/// a constant. Measured with four DRAM-bandwidth hogs on a 16-core/32-thread +/// host: compact 4.54 ms/token against spread 5.03-6.26. With eight hogs +/// covering both halves the ranking inverts (compact 15.92 against spread +/// 6.08). See [`crate::core_topology`]'s module docs for the full table. +pub const DECODE_PLACEMENT_ENV: &str = "ONNX_GENAI_CPU_DECODE_PLACEMENT"; + +/// The complete set of accepted placement modes. +const ACCEPTED_PLACEMENTS: &str = "`spread`, `compact`"; + +/// How a decode pool's workers are laid out across the CPUs it may use. +/// +/// This is a *policy*, and the point of naming it is that it can be selected +/// rather than assumed. A test that asserts one worker per physical core is +/// asserting `Spread`, not correctness, and must say so -- see #1802. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub enum CorePlacement { + /// One worker per physical core for as long as cores last. + #[default] + Spread, + /// SMT siblings filled before the next core is used. + Compact, +} + +impl CorePlacement { + /// Parse a [`DECODE_PLACEMENT_ENV`] value. Absent or empty is the default. + pub fn parse(raw: Option<&str>) -> std::result::Result { + match raw.map(str::trim) { + None | Some("") => Ok(Self::default()), + Some("spread") => Ok(Self::Spread), + Some("compact") => Ok(Self::Compact), + Some(other) => Err(format!( + "{DECODE_PLACEMENT_ENV}=`{other}` is not a recognized placement mode; \ + accepted modes are {ACCEPTED_PLACEMENTS}" + )), + } + } + + /// The placement selected for this process, read once. + /// + /// An unparseable value is reported once and then treated as the default. + /// The sibling [`DecodeAffinity`] selector rejects hard because it decides + /// *which* CPUs the process may touch and a wrong answer there is a + /// correctness question; this one decides only how workers are arranged + /// inside a set already chosen, so refusing to decode because a tuning knob + /// is misspelled would be the worse failure. It is still never silent. + /// + /// `OnceLock` makes "reported once" structural rather than a promise: the + /// closure runs at most once per process. + pub fn from_env() -> Self { + static V: OnceLock = OnceLock::new(); + *V.get_or_init(|| Self::from_value(std::env::var(DECODE_PLACEMENT_ENV).ok().as_deref())) + } + + /// [`Self::from_env`] without the cache or the environment read. + /// + /// Split out because the fallback is otherwise untestable: `from_env` + /// memoizes process-wide, so a test driving it would fix the policy for + /// every other test in the binary and could still only ever observe one + /// value. This half is the part with a decision in it. + fn from_value(raw: Option<&str>) -> Self { + match Self::parse(raw) { + Ok(placement) => placement, + Err(message) => { + eprintln!( + "onnx-genai: {message}; using `{}`", + Self::default().as_str() + ); + Self::default() + } + } + } + + /// The wire/report name, so a diagnostic and a test read the same string. + pub fn as_str(self) -> &'static str { + match self { + Self::Spread => "spread", + Self::Compact => "compact", + } + } +} + /// Render the discovered available-node list for a diagnostic, or state plainly /// that topology is unavailable, so every invalid value reports the same three /// facts: what was rejected, what is accepted, and what nodes actually exist. @@ -534,6 +629,21 @@ pub fn set_current_thread_affinity(cpus: &[usize]) -> std::result::Result<(), St Err("process-wide CPU affinity masking is only implemented on Linux (no-op)".to_string()) } +/// Whether [`set_current_thread_affinity`] can do anything on this target. +/// +/// A test that narrows its own cpuset to build a controlled shape has to tell +/// two failures apart: "this target has no process-wide masking, so the shape is +/// unbuildable by construction" and "masking exists here and did not work". +/// The first is a legitimate skip; the second is a defect, and reading it as the +/// first is how a whole platform's coverage disappears without anyone noticing. +/// +/// Derived from `cfg!` and *not* from a runtime call, for the same reason +/// [`crate::core_topology::DETECTION_SUPPORTED`] is: a support flag computed by +/// asking the thing whether it worked answers "unsupported" for every runtime +/// failure, which is exactly the collapse it exists to prevent. +#[cfg(test)] +pub(crate) const AFFINITY_MASKING_SUPPORTED: bool = cfg!(target_os = "linux"); + /// Whether the user explicitly requested a decode affinity via /// [`DECODE_AFFINITY_ENV`] (a non-empty value). When set, the automatic /// good-citizen process mask stands down so the user's affinity choice wins. @@ -979,15 +1089,36 @@ pub fn plan_decode_affinity(worker_count: usize) -> std::result::Result, +) -> Vec { + order_pin_targets_for(cpus, cores, CorePlacement::from_env()) +} + +/// Order a decode pool's CPU list for an explicitly named placement policy. /// /// [`crate::kernels::matmul_nbits`]'s `build_decode_pool` pins worker `i` to -/// `cpus[i % cpus.len()]`, so the *order* of this list — not just its contents — -/// decides whether an 8-worker pool on a 16-core/32-thread node occupies 8 cores -/// or 4. Same set, same cpuset guarantees; only the ranking changes. With no -/// discoverable SMT map this is the identity. +/// `cpus[i % cpus.len()]`, so the *order* of this list — not just its contents +/// — decides whether an 8-worker pool on a 16-core/32-thread node occupies 8 +/// cores or 4. Same set, same cpuset guarantees; only the ranking changes. With +/// no discoverable SMT map both policies are the identity, because there is +/// nothing to rank by. +/// +/// The returned list is always a permutation of `cpus`: a CPU the topology does +/// not know about keeps its place at the end rather than being dropped. The +/// caller's set is a *permission*, and silently shrinking it would shrink the +/// pool. (A *repeated* CPU is de-duplicated rather than repeated, under both +/// policies and both before and after #1805's follow-up; the callers are +/// cpusets, which cannot contain one.) /// /// # Who uses it, and a reversal /// @@ -1013,18 +1144,39 @@ pub fn plan_decode_affinity(worker_count: usize) -> std::result::Result, + placement: CorePlacement, ) -> Vec { let Some(cores) = cores else { return cpus.to_vec(); }; - let leaders = cores.leaders_within(cpus); - let leader_set: std::collections::BTreeSet = leaders.iter().copied().collect(); - let mut ordered = leaders; - ordered.extend(cpus.iter().copied().filter(|cpu| !leader_set.contains(cpu))); - ordered + match placement { + CorePlacement::Spread => { + let leaders = cores.leaders_within(cpus); + let leader_set: std::collections::BTreeSet = leaders.iter().copied().collect(); + let mut ordered = leaders; + ordered.extend(cpus.iter().copied().filter(|cpu| !leader_set.contains(cpu))); + ordered + } + CorePlacement::Compact => { + let allowed: std::collections::BTreeSet = cpus.iter().copied().collect(); + let mut ordered = Vec::with_capacity(cpus.len()); + let mut placed: std::collections::BTreeSet = std::collections::BTreeSet::new(); + for leader in cores.leaders_within(cpus) { + for sibling in cores.siblings_of(leader).unwrap_or(&[]) { + if allowed.contains(sibling) && placed.insert(*sibling) { + ordered.push(*sibling); + } + } + } + // Whatever the topology could not account for keeps its original + // relative order at the end, so this stays a permutation. + ordered.extend(cpus.iter().copied().filter(|cpu| !placed.contains(cpu))); + ordered + } + } } /// Byte-level layout of the records `GetLogicalProcessorInformationEx` writes, @@ -1322,6 +1474,38 @@ mod windows_imp { mod tests { use super::*; + /// The support flag must not claim less than the call delivers. + /// + /// The drift that matters is one-directional. If someone implements + /// process-wide masking on Windows and leaves + /// [`AFFINITY_MASKING_SUPPORTED`] at `cfg!(target_os = "linux")`, every + /// caller that uses the flag to decide "unbuildable by construction, skip" + /// keeps skipping on a platform that could now answer -- silently, forever. + /// That is the fail-open direction and it is what this asserts: a call that + /// *succeeded* proves support, whatever the flag says. + /// + /// The converse -- flag `true`, call failed -- is deliberately **not** + /// asserted here. It would false-fail on a CPU hot-unplugged between + /// reading the mask and re-applying it, and its real failure mode is + /// already loud: every caller of this flag panics rather than skips when + /// masking is supported and does not work. + #[test] + fn a_successful_affinity_call_is_never_reported_as_an_unsupported_target() { + let Some(allowed) = allowed_cpus() else { + return; + }; + // Re-applying the mask already in force: a no-op on success, and it + // cannot widen or narrow the thread running this test either way. + let succeeded = set_current_thread_affinity(&allowed).is_ok(); + assert!( + !succeeded || AFFINITY_MASKING_SUPPORTED, + "set_current_thread_affinity succeeded on a target where \ + AFFINITY_MASKING_SUPPORTED is false. Callers read that flag as \ + `this shape is unbuildable here, skip`, so the flag is now \ + disabling checks on a platform that can run them." + ); + } + /// Build one `RelationNumaNode` record exactly as the OS lays it out: an /// 8-byte header, then `NUMA_NODE_RELATIONSHIP`, then the node's /// `GROUP_AFFINITY`. `size` overrides the declared length so a test can @@ -2171,10 +2355,149 @@ mod tests { fn pin_targets_are_ordered_one_per_core_then_siblings() { let smt = adjacent_smt_topology(); let node: Vec = (0..8).collect(); + // Names the policy it is asserting. `order_pin_targets` reads the + // process-wide selector, so asserting the spread ordering through the + // wrapper would silently be asserting "this process happens to be + // spread" -- the exact conflation #1802 asks to stop making. + let ordered = order_pin_targets_for(&node, Some(&smt), CorePlacement::Spread); // The pool builder pins worker `i` to `cpus[i % len]`, so the first four // workers of an 8-CPU node must land on four different cores. - let ordered = order_pin_targets(&node, Some(&smt)); assert_eq!(ordered, vec![0, 2, 4, 6, 1, 3, 5, 7]); - assert_eq!(order_pin_targets(&node, None), node); + assert_eq!( + order_pin_targets_for(&node, None, CorePlacement::Spread), + node + ); + // ...and the default really is the policy this asserts, stated + // separately so a changed default fails with its own name rather than + // looking like an ordering bug. + assert_eq!(CorePlacement::default(), CorePlacement::Spread); + } + + #[test] + fn placement_parse_accepts_exactly_the_documented_modes() { + assert_eq!(CorePlacement::parse(None).unwrap(), CorePlacement::Spread); + assert_eq!( + CorePlacement::parse(Some("")).unwrap(), + CorePlacement::Spread + ); + assert_eq!( + CorePlacement::parse(Some(" ")).unwrap(), + CorePlacement::Spread + ); + assert_eq!( + CorePlacement::parse(Some(" spread ")).unwrap(), + CorePlacement::Spread + ); + assert_eq!( + CorePlacement::parse(Some("compact")).unwrap(), + CorePlacement::Compact + ); + // Every accepted mode round-trips through the name the report and the + // diagnostics use, so a test and a log line cannot drift apart. + for placement in [CorePlacement::Spread, CorePlacement::Compact] { + assert_eq!( + CorePlacement::parse(Some(placement.as_str())).unwrap(), + placement + ); + } + // Rejected, and the rejection has to say what was wrong and what is + // accepted -- an unhelpful error is how a misspelled knob becomes a + // silent default. + for bad in ["Spread", "COMPACT", "sprd", "1", "true", "one-per-core"] { + let message = CorePlacement::parse(Some(bad)) + .expect_err(&format!("`{bad}` must not be accepted as a placement")); + assert!( + message.contains(bad) && message.contains(DECODE_PLACEMENT_ENV), + "the rejection must name the variable and the offending value: {message}" + ); + assert!( + message.contains("spread") && message.contains("compact"), + "the rejection must list the accepted modes: {message}" + ); + } + } + + #[test] + fn an_unparseable_placement_falls_back_without_going_silent() { + // The fallback path, which `from_env` alone cannot exercise: it caches + // process-wide, so whichever value the first caller in the binary sees + // is the only one any test could ever observe. + for bad in ["Spread", "one-per-core", "1"] { + assert_eq!( + CorePlacement::from_value(Some(bad)), + CorePlacement::default(), + "an unrecognized placement must fall back to the default rather than \ + selecting something arbitrary" + ); + } + // ...and the accepted values are not merely tolerated by the same path. + assert_eq!( + CorePlacement::from_value(Some("compact")), + CorePlacement::Compact, + "the fallback path swallowed a value it was supposed to honor, which would \ + make the knob inert while looking like it parsed" + ); + assert_eq!(CorePlacement::from_value(None), CorePlacement::default()); + } + + #[test] + fn compact_and_spread_are_permutations_that_differ_on_an_smt_host() { + let smt = adjacent_smt_topology(); + let node: Vec = (0..8).collect(); + let spread = order_pin_targets_for(&node, Some(&smt), CorePlacement::Spread); + let compact = order_pin_targets_for(&node, Some(&smt), CorePlacement::Compact); + + // A permission set that shrinks would shrink the pool, so both policies + // must be permutations of their input and not merely subsets. + for (name, ordered) in [("spread", &spread), ("compact", &compact)] { + let mut sorted = (*ordered).clone(); + sorted.sort_unstable(); + assert_eq!(sorted, node, "`{name}` is not a permutation of its input"); + } + + assert_ne!( + spread, compact, + "the two policies produced the same ranking on an SMT host, so the \ + selector parses and changes nothing" + ); + assert_eq!(compact, vec![0, 1, 2, 3, 4, 5, 6, 7]); + + // The property the ordering exists for: `build_decode_pool` pins worker + // `i` to `cpus[i % len]`, so the first four entries decide how many + // cores a four-worker pool occupies. + assert_eq!(smt.physical_cores_within(&spread[..4]), 4); + assert_eq!(smt.physical_cores_within(&compact[..4]), 2); + } + + #[test] + fn a_cpu_the_topology_does_not_know_keeps_its_place_under_either_policy() { + let smt = adjacent_smt_topology(); + // 99 is outside the sibling map; 16/17 are a known-shaped pair that the + // 8-core map above also does not cover. + let node = vec![0, 1, 99, 4]; + for placement in [CorePlacement::Spread, CorePlacement::Compact] { + let ordered = order_pin_targets_for(&node, Some(&smt), placement); + let mut sorted = ordered.clone(); + sorted.sort_unstable(); + assert_eq!( + sorted, + vec![0, 1, 4, 99], + "`{}` dropped or duplicated a CPU it did not recognize", + placement.as_str() + ); + } + } + + #[test] + fn without_a_topology_every_policy_is_the_identity() { + let node = vec![7, 3, 11, 0]; + for placement in [CorePlacement::Spread, CorePlacement::Compact] { + assert_eq!( + order_pin_targets_for(&node, None, placement), + node, + "`{}` reordered a list it has no ranking information for", + placement.as_str() + ); + } } } diff --git a/crates/onnx-runtime-ep-cpu/src/decode_spmd.rs b/crates/onnx-runtime-ep-cpu/src/decode_spmd.rs index ecf5a055db..707d448631 100644 --- a/crates/onnx-runtime-ep-cpu/src/decode_spmd.rs +++ b/crates/onnx-runtime-ep-cpu/src/decode_spmd.rs @@ -3917,10 +3917,13 @@ fn explicit_affinity_shards_for( if cpus.is_empty() { return None; } - // Order through the same spread the default path uses. An explicit - // request selects *which* CPUs, never how workers are laid out - // across them; without this, `compact` would pin two workers per - // physical core and quietly reintroduce the defect #1729 fixed. + // Order through the same placement policy the default path uses. An + // explicit request selects *which* CPUs, never how workers are laid + // out across them; without this, `compact` (a NUMA-node selector) + // would inherit whatever ordering happened to be lying around and + // could quietly reintroduce the defect #1729 fixed. The layout + // itself is chosen by `ONNX_GENAI_CPU_DECODE_PLACEMENT`, which + // defaults to spread. let cores = crate::core_topology::host(); let cpus = crate::decode_affinity::order_pin_targets(&cpus, cores); let core_count = cores.map_or(0, |cores| cores.leaders_within(&cpus).len()); @@ -4691,7 +4694,11 @@ mod tests { (0..8).map(|c| vec![c * 2, c * 2 + 1]), ); let all: Vec = (0..16).collect(); - let ordered = crate::decode_affinity::order_pin_targets(&all, Some(&synthetic)); + let ordered = crate::decode_affinity::order_pin_targets_for( + &all, + Some(&synthetic), + crate::decode_affinity::CorePlacement::Spread, + ); let mut leaders = synthetic.leaders_within(&all); leaders.sort_unstable(); let mut first_eight: Vec = ordered.iter().take(8).copied().collect(); @@ -6112,9 +6119,10 @@ mod tests { } // Two full physical cores' worth of CPUs, spread the way the placement // policy spreads them. - let cpus: Vec = crate::decode_affinity::order_pin_targets( + let cpus: Vec = crate::decode_affinity::order_pin_targets_for( &(0..cores.logical_count()).collect::>(), Some(cores), + crate::decode_affinity::CorePlacement::Spread, ); let workers = cores.core_count().min(4); if !environment_can_pin(&cpus[..workers.min(cpus.len())]) { @@ -6221,9 +6229,10 @@ mod tests { ); return; } - let cpus: Vec = crate::decode_affinity::order_pin_targets( + let cpus: Vec = crate::decode_affinity::order_pin_targets_for( &(0..cores.logical_count()).collect::>(), Some(cores), + crate::decode_affinity::CorePlacement::Spread, ); let workers = cores.core_count().min(4); if !environment_can_pin(&cpus[..workers.min(cpus.len())]) { @@ -6421,9 +6430,10 @@ mod tests { ) { return; } - let cpus: Vec = crate::decode_affinity::order_pin_targets( + let cpus: Vec = crate::decode_affinity::order_pin_targets_for( &(0..cores.logical_count()).collect::>(), Some(cores), + crate::decode_affinity::CorePlacement::Spread, ); let workers = cores.core_count().min(4).min(cpus.len()); if workers < 2 { @@ -6948,9 +6958,10 @@ mod tests { ) { return; } - let spread = crate::decode_affinity::order_pin_targets( + let spread = crate::decode_affinity::order_pin_targets_for( &(0..cores.logical_count()).collect::>(), Some(cores), + crate::decode_affinity::CorePlacement::Spread, ); if !environment_can_pin(&spread[..2]) { return; diff --git a/crates/onnx-runtime-ep-cpu/src/kernels/matmul_nbits.rs b/crates/onnx-runtime-ep-cpu/src/kernels/matmul_nbits.rs index 8424081a66..75d508f9f0 100644 --- a/crates/onnx-runtime-ep-cpu/src/kernels/matmul_nbits.rs +++ b/crates/onnx-runtime-ep-cpu/src/kernels/matmul_nbits.rs @@ -20026,6 +20026,18 @@ mod tests { /// `Some(true)` on a build whose pinning reports success and does /// nothing. `realized` is the one that asks the kernel. planned_distinct_cores: Option, + /// Which placement policy the child selected, as + /// [`crate::decode_affinity::CorePlacement::as_str`]. + /// + /// Reported so the sweep can say what it is measuring instead of + /// assuming it. An assertion about *where* workers run is an assertion + /// about this value, and one that does not name it is asserting a + /// policy it never selected. + placement_policy: String, + /// The cpuset the child actually measured, so the parent can predict a + /// layout for *that* set rather than for whatever it can see itself. + /// Empty when the child could not read its own affinity. + allowed_cpus: Vec, /// What the pool's workers *observed* about their own placement, as the /// wire code of [`crate::decode_spmd::RealizedPlacement`]. The blind-spot /// codes are distinct from `one-per-core`, so an unanswerable child can @@ -20051,6 +20063,17 @@ mod tests { /// old name would have left the parent's escape hatch below checking a /// host fact that no longer has anything to do with it. worker_masks_readable: bool, + /// Whether the child successfully narrowed itself to one CPU per + /// physical core. Only the default arm asks for this; every other arm + /// reports `false` because it never tried. + /// + /// A separate field rather than an inference from `allowed == cores`, + /// because that equality is also what the arm *asserts*: deriving the + /// precondition from the conclusion is how a check comes to confirm + /// itself. On a target with no process-wide masking the two are + /// genuinely different facts -- narrowing did not happen, and the + /// equality may still hold on a machine without SMT. + narrowed_to_leaders: bool, /// Attempts the parent consumed to obtain this report. `1` on a clean /// run; more means environmental crashes were retried away and this /// number is the only in-band trace of them (#1745). @@ -20312,20 +20335,35 @@ mod tests { // `requested=0` marks the default-resolver arm in the report line: the // parent did not ask for a width, so there is no requested value to // echo and the only honest thing to print is "none". + let mut narrowed = false; let requested: usize = if raw == SPMD_WIDTH_CHILD_DEFAULT { - restrict_self_to_leader_cpus(); + narrowed = restrict_self_to_leader_cpus(); 0 } else { raw.parse().expect("the parent passes a decimal width") }; - let allowed = crate::decode_affinity::allowed_cpus().map_or(0, |cpus| cpus.len()); + // Read *after* any narrowing above, so the report describes the cpuset + // the pool was actually built on rather than the one inherited. + let allowed_cpus = crate::decode_affinity::allowed_cpus().unwrap_or_default(); + let allowed = allowed_cpus.len(); + // The set, not just its size. The parent predicts the layout a policy + // produces by ordering *these* CPUs, so an equal-sized but differently + // populated cpuset would leave it grading another machine. + let cpulist = if allowed_cpus.is_empty() { + "na".to_string() + } else { + allowed_cpus + .iter() + .map(usize::to_string) + .collect::>() + .join(",") + }; // Fails closed on a supported target: silently substituting the logical // count for the physical one here would hand the parent an inflated // `cores` and make its placement expectation trivially satisfiable. - let cores = - crate::core_topology::require_host_for_placement().map_or(allowed, |topology| { - crate::decode_affinity::allowed_cpus() - .map_or(allowed, |cpus| topology.physical_cores_within(&cpus)) + let cores = crate::core_topology::require_host_for_placement() + .map_or(allowed, |topology| { + topology.physical_cores_within(&allowed_cpus) }); let (pool_built, workers, nodes, placement, honest, fully_pinned, realized, observed_ok) = match crate::decode_spmd::pools() { @@ -20360,13 +20398,15 @@ mod tests { println!( "{SPMD_WIDTH_MARKER}requested={requested} available={} allowed={allowed} \ nodes={nodes} workers={workers} pool={} cores={cores} placement={} honest={} \ - pinned={} proc={} realized={realized}", + pinned={} proc={} realized={realized} policy={} narrowed={} cpulist={cpulist}", available_parallelism(), u8::from(pool_built), tri(placement), tri(honest), u8::from(fully_pinned), - u8::from(observed_ok) + u8::from(observed_ok), + crate::decode_affinity::CorePlacement::from_env().as_str(), + u8::from(narrowed) ); quiesce_pools_for_child_exit("spmd_realized_width"); } @@ -20380,51 +20420,116 @@ mod tests { /// free of a `taskset` dependency and makes the mask an observable the /// child reports (`allowed`/`cores`) rather than an assumption. /// - /// Only Linux can do this. `set_current_thread_affinity` is a documented - /// no-op that returns `Err` everywhere else (see `decode_affinity.rs`), so - /// on those targets the leader-only cpuset cannot be constructed at all and - /// the caller must skip rather than assert. That is a different thing from - /// the restriction failing, which stays fatal on Linux. - fn restrict_self_to_leader_cpus() { - if !cfg!(target_os = "linux") { - return; + /// Returns whether the narrowing actually happened, which the parent needs + /// in order to tell two different "no" answers apart: + /// + /// * **Unbuildable by construction.** `set_current_thread_affinity` is + /// implemented on Linux only, so on Windows and macOS this shape cannot be + /// created at all and the arm has nothing to measure. That is a skip, and + /// it is stated on stderr rather than passed silently. + /// * **Supported and broken.** On a target that *has* masking, a failure + /// leaves the child on the full cpuset, where the arm's `workers == cores` + /// would be a claim about the host rather than about the resolver. That + /// panics. + /// + /// The original version panicked on both, which is why #2059 merged with + /// both Windows lanes red: it read "this platform has no process-wide + /// masking" as a defect. Collapsing them the other way -- skipping on both + /// -- would have been worse, because it removes Linux coverage the moment + /// masking breaks there, and Linux is where this arm actually runs. + /// + /// #2078 fixed the red with the platform half of this, which is the part + /// that matters for a green lane. Reporting the *outcome* rather than the + /// platform is what closes the remainder: the parent then keys its skip on + /// what happened rather than on what it assumes must have happened, and the + /// three unanswerable paths below stop being silent `return`s. + fn restrict_self_to_leader_cpus() -> bool { + // Short-circuit before touching topology or cpuset state. On a target + // with no masking the answer is already settled, and probing on the way + // to a foregone conclusion only creates ways for this arm to fail for + // reasons that are not its subject -- `require_host_for_placement` + // below fails closed on Windows, which is right in general and wrong + // as a side effect of a call that was never going to narrow anything. + if !crate::decode_affinity::AFFINITY_MASKING_SUPPORTED { + eprintln!( + "leader-only narrowing skipped: process-wide CPU affinity masking is \ + implemented only on Linux, so a leader-only cpuset cannot be constructed \ + on this target" + ); + return false; } let Some(allowed) = crate::decode_affinity::allowed_cpus() else { - return; + eprintln!( + "leader-only narrowing skipped: this target cannot report the CPUs it is \ + allowed to run on, so there is no set to narrow" + ); + return false; }; - let Ok(topology) = crate::core_topology::require_host_for_placement() else { - return; + let topology = match crate::core_topology::require_host_for_placement() { + Ok(topology) => topology, + Err(reason) => { + eprintln!("leader-only narrowing skipped: {reason}"); + return false; + } }; let leaders = topology.leaders_within(&allowed); if leaders.is_empty() { - return; + eprintln!( + "leader-only narrowing skipped: the detected topology covers none of the \ + {} allowed CPUs, so it names no leaders to narrow to", + allowed.len() + ); + return false; + } + match crate::decode_affinity::set_current_thread_affinity(&leaders) { + Ok(()) => true, + // Fails closed exactly where the capability exists. `cfg!`-derived, + // so it cannot quietly become `false` on a host that should have + // answered -- the same reason `DETECTION_SUPPORTED` is not a + // runtime probe. + Err(reason) if crate::decode_affinity::AFFINITY_MASKING_SUPPORTED => panic!( + "restrict the child to leader CPUs: {reason}. This target implements \ + process-wide affinity masking, so this is a failure rather than an \ + absent capability, and continuing would leave the child on the full \ + cpuset where `workers == cores` holds for the wrong reason." + ), + Err(reason) => { + eprintln!("leader-only narrowing skipped: {reason}"); + false + } } - // A failure here is not silently tolerated: it would leave the child on - // the full cpuset, where `workers == cores` holds for the wrong reason - // and the assertion would pass without testing anything. - crate::decode_affinity::set_current_thread_affinity(&leaders) - .expect("restrict the child to leader CPUs"); } fn realized_width_report(requested: usize) -> RealizedWidth { - realized_width_report_with(requested, false) + realized_width_report_with(requested, false, None) } /// Spawn the child with **no** explicit width, on a leader-only cpuset, so /// the default resolver is the thing under test. fn realized_default_width_report() -> RealizedWidth { - realized_width_report_inner(SPMD_WIDTH_CHILD_DEFAULT, None, false) + realized_width_report_inner(SPMD_WIDTH_CHILD_DEFAULT, None, false, None) } - fn realized_width_report_with(requested: usize, dishonest: bool) -> RealizedWidth { + /// The sweep's child, with the placement policy named rather than inherited. + /// + /// `placement: None` removes [`crate::decode_affinity::DECODE_PLACEMENT_ENV`] + /// entirely, which is what "the default" means here -- the point of the + /// default arm is to measure whatever policy ships, not a policy the test + /// pinned down for its own convenience. + fn realized_width_report_with( + requested: usize, + dishonest: bool, + placement: Option, + ) -> RealizedWidth { let width = requested.to_string(); - realized_width_report_inner(&width, Some(&width), dishonest) + realized_width_report_inner(&width, Some(&width), dishonest, placement) } fn realized_width_report_inner( child_value: &str, explicit_width: Option<&str>, dishonest: bool, + placement: Option, ) -> RealizedWidth { let mut command = std::process::Command::new(std::env::current_exe().unwrap()); command @@ -20461,6 +20566,17 @@ mod tests { .env_remove("RAYON_NUM_THREADS"); } } + match placement { + Some(policy) => { + command.env( + crate::decode_affinity::DECODE_PLACEMENT_ENV, + policy.as_str(), + ); + } + None => { + command.env_remove(crate::decode_affinity::DECODE_PLACEMENT_ENV); + } + } // Never inherited: the sweep's honesty assertion would be meaningless // if an ambient value could switch the injection on, and unfalsifiable // if one could switch it off. @@ -20553,9 +20669,21 @@ mod tests { cores: field("cores"), planned_distinct_cores: tri("placement"), realized: text("realized").to_string(), + placement_policy: text("policy").to_string(), + allowed_cpus: match text("cpulist") { + "na" => Vec::new(), + list => list + .split(',') + .map(|cpu| { + cpu.parse() + .unwrap_or_else(|_| panic!("`cpulist` is not a CPU list: {report}")) + }) + .collect(), + }, placement_honest: tri("honest"), fully_pinned: field("pinned") == 1, worker_masks_readable: field("proc") == 1, + narrowed_to_leaders: field("narrowed") == 1, attempts: attempts_used, }; @@ -20646,8 +20774,8 @@ mod tests { if allowed < 2 { return; } - let honest = realized_width_report_with(2, false); - let lying = realized_width_report_with(2, true); + let honest = realized_width_report_with(2, false, None); + let lying = realized_width_report_with(2, true, None); // A pool that did not fully pin cannot be cross-checked at all, so the // experiment below has no signal to read either way. Skip on the *host @@ -20743,7 +20871,17 @@ mod tests { // `--ignored`: it must still not assert on an unrestricted cpuset. Its // `eprintln!` is only rendered under `--nocapture`, which is why it is // not the primary signal. - if !cfg!(target_os = "linux") { + // + // Keyed on the same `cfg!`-derived constant the child uses, rather than + // on a third copy of `cfg!(target_os = "linux")`, so that the day + // Windows masking is implemented there is one place to change and a + // unit test (`a_successful_affinity_call_is_never_reported_as_an_unsupported_target`) + // that fails if it is not changed. The attribute + // above is unavoidably a second copy -- `ignore` takes a `cfg`, not a + // `const` -- but its drift failure mode is `ignored`, which is loud and + // is not counted as executed. This arm's would be a silent `ok`, so it + // is the one that must not have its own copy. + if !crate::decode_affinity::AFFINITY_MASKING_SUPPORTED { eprintln!( "SKIP a_default_width_pool_on_leader_cpus_uses_every_core_it_was_given: \ process-wide CPU affinity masking is implemented only on Linux, so a \ @@ -20754,6 +20892,24 @@ mod tests { let report = realized_default_width_report(); + // Past the platform gate, narrowing is not optional. This asserts what + // *happened* rather than re-deriving what the platform should have + // allowed: the child has three further paths on which it declines to + // narrow (no readable cpuset, undetectable topology, a topology + // covering none of the allowed CPUs), and every one of them used to be + // a silent `return` that left this arm measuring the full cpuset. + // + // `allowed == cores` below catches that only on an SMT host, where the + // two differ. On a machine without SMT it holds on the *un-narrowed* + // set too, and the arm would pass while testing nothing. + assert!( + report.narrowed_to_leaders, + "this target implements process-wide affinity masking, so the child must have \ + narrowed itself to one CPU per physical core -- it reports that it did not, \ + and every assertion below would then be about the host rather than about the \ + default width resolver ({report:?})" + ); + if !report.pool_built { // No persistent pool on this target: there is no placement to // check, and inventing a pass here would be the vacuous case this @@ -20793,6 +20949,24 @@ mod tests { } ); + // Deliberately kept as an unconditional claim here, even though the + // rest of this sweep no longer asserts one-per-core under the default + // policy (see #1802 and the long note in + // `every_benchmarked_decode_width_realizes_the_worker_count_it_requests`). + // It is not a policy contract *in this arm*: the child narrowed itself + // to one CPU per physical core, `order_pin_targets_for` is a + // permutation of the allowed CPUs under every policy it implements -- + // `Compact` dedups through `placed` and appends the remainder exactly + // once -- so on a leader-only cpuset compact and spread both land one + // worker per core. The mask shape forces the outcome, not the policy. + // + // The assumption that buys that, stated so it fails loudly rather than + // silently if it stops holding: the default must place *at most one + // worker per allowed CPU*. An oversubscribing shared-core default would + // realize `shared-core` here (and trip `workers == cores` above), and + // the right response then is to gate this arm on the policy, not to + // delete it -- the width claim it exists to make (#1780) is + // policy-neutral and still worth asserting. assert_eq!( report.realized, "one-per-core", "a default-width pool on a leader-only cpuset must realize one \ @@ -20950,75 +21124,79 @@ mod tests { report.worker_masks_readable ); - // (2) POLICY-DEPENDENT, and deliberately labelled as such: one - // worker per physical core is #1729's *spread* policy, not a law. - // It is measured better on a quiet host and measured ~26% worse - // with a single ~90% co-tenant, because a compact pool leaves half - // the box for the co-tenant to land on. The standing user direction + // (2) The pool must report *which* placement policy it applied, + // and that report must be one the child could actually have + // produced. This replaces an unconditional one-worker-per-physical + // -core assertion that used to live here. + // + // Why it had to go: one worker per physical core is #1729's + // *spread* policy, not a law. It measures better on a quiet host + // and ~26% worse with a single ~90% co-tenant, because a compact + // pool leaves half the box for the co-tenant to land on. The + // standing user direction // (`.squad/decisions/inbox/copilot-cpu-shared-host-default-2026-08-23.md`, // via #1729) is that a policy which wins only under exclusive - // quiet-host conditions is not a valid default, so this assertion - // may legitimately have to change. + // quiet-host conditions is not a valid default -- so the default + // sweep asserting that policy makes #1802's shared default + // unlandable, and the path of least resistance for whoever hits it + // is to weaken the assertion rather than fix its shape. // - // If it does, change it *deliberately* -- it is an assertion about - // the policy that is currently shipped, not about correctness. Do - // not "fix" it by loosening it to match whatever the pool did, - // which would delete the check; and do not read a failure here as a - // regression without first deciding which policy is intended. Part - // (1) above is the part that must survive that decision intact. - // - // Only assert where one-per-core is achievable: past the core - // budget the workers must double up, and demanding otherwise would - // fail the host, not the code. - if report.workers <= report.cores && report.planned_distinct_cores.is_some() { - assert_eq!( - report.planned_distinct_cores, - Some(true), - "requested width {requested} realized {} workers that share physical \ - cores on a host with {} cores available, so a decode row labelled \ - t={requested} ran on fewer front ends than its label implies. This is \ - an assertion about the currently shipped spread policy ({report:?})", - report.workers, - report.cores - ); - } - saw_placement_check |= - report.workers <= report.cores && report.planned_distinct_cores.is_some(); - - // (3) The same policy claim, asked of the kernel instead of the - // planner. (2) reads `worker_cpus()`, which is the assignment plus - // whatever the pin call *returned*, so it survives a build whose - // pinning does nothing and reports success. This one reads the mask - // each worker read for itself. - // - // Blind spots are named, and none of them is `one-per-core`: a - // child that could not answer fails the guard below rather than - // passing this one. Conditioned on the affinity query existing -- - // `fully_pinned` reports the *pinning* capability, and a target - // that can pin without being able to read a mask back would fail - // this for a host limitation rather than a defect. The guard - // immediately after is what stops that condition from becoming an - // escape: where the query does exist, a child claiming it does not - // is a defect in the apparatus. - if crate::decode_affinity::affinity_observation_supported() - && report.fully_pinned - && report.workers <= report.cores - { - assert_eq!( - report.realized, "one-per-core", - "requested width {requested}: every worker reports an applied pin and the host has room, but the placement the workers actually observed for themselves is `{}` -- so the label t={requested} describes a placement the kernel did not enforce ({report:?})", - report.realized - ); - } + // The honest, policy-neutral claim -- "the pool places workers + // where it says it places them" -- is part (1) above, and it holds + // under every policy. Where the *spread* policy is what is being + // asked for, it is asserted explicitly, by + // `an_explicit_spread_policy_places_one_worker_per_physical_core`, + // which selects it rather than assuming it. + // Equality with the shipped default rather than membership in the + // accepted set. Membership could not fail: the child prints + // `CorePlacement::as_str()`, which is total over the enum, so the + // field is always one of them. Equality is a real claim -- this arm + // removes `DECODE_PLACEMENT_ENV` from the child, so it fails if an + // ambient value leaks in or `from_env` reads the wrong variable, + // and it has to be updated deliberately when #1802 flips the + // default, which is the discipline this test is arguing for. + assert_eq!( + report.placement_policy, + crate::decode_affinity::CorePlacement::default().as_str(), + "width {requested}: the default arm removes the placement variable from \ + the child, so the child must report the policy that ships -- it reported \ + `{}` ({report:?})", + report.placement_policy + ); + + // The realized placement must be a code the apparatus can act on. + // `no-affinity-query` on a target that *has* the query means the + // observation was switched off, which is the exact shape of a check + // that cannot fail. Note this asserts nothing about *where* the + // workers went. assert!( !crate::decode_affinity::affinity_observation_supported() || report.realized != "no-affinity-query", - "width {requested}: this target has a per-thread affinity query, yet the child reported it had none -- the realized-placement observation is switched off, which is the exact shape of a check that cannot fail ({report:?})" + "width {requested}: this target has a per-thread affinity query, yet the \ + child reported it had none -- the realized-placement observation is \ + switched off, which is the exact shape of a check that cannot fail \ + ({report:?})" + ); + assert!( + report.realized != "affinity-query-failed", + "width {requested}: at least one worker could not read the affinity in \ + force for it, so its placement is unknown rather than correct -- an \ + unreadable mask must never be scored as a placed worker ({report:?})" ); + + // Anti-vacuity for the honesty half, per width. `saw_placement_check` + // used to count widths at which the *planner's* one-per-core policy + // was asserted; there is no such assertion here any more, so it + // counts widths at which a policy-neutral verdict was actually + // produced. A counter that outlives the check it counted is worse + // than no counter, because it keeps reporting coverage that stopped + // existing. + saw_placement_check |= report.placement_honest.is_some(); saw_realized_placement_check |= crate::decode_affinity::affinity_observation_supported() && report.fully_pinned - && report.workers <= report.cores; + && report.workers <= report.cores + && report.placement_honest == Some(true); assert_eq!( report.workers, effective, @@ -21064,16 +21242,21 @@ mod tests { } if allowed_now >= 2 && crate::core_topology::DETECTION_SUPPORTED { assert!( - saw_placement_check, - "no swept width produced a checkable *planned* placement on a \ - {allowed_now}-CPU cpuset, so the planner's core layout went unasserted \ - and only the worker count was verified" + saw_placement_check + || !crate::decode_affinity::pinning_supported() + || !crate::decode_affinity::affinity_observation_supported(), + "no swept width produced a placement-honesty verdict on a \ + {allowed_now}-CPU cpuset that supports pinning and can read affinity \ + back, so `worker_cpus()` went unverified at every width and the sweep \ + checked the worker count only" ); // Separate counter from `saw_placement_check`, because the two can - // diverge in exactly the direction that matters: the planned check - // needs only a topology, while the realized one needs pins that - // actually took. A single flag would let the planner's assertions - // vouch for a realized check that never ran. + // diverge in exactly the direction that matters: a verdict of + // `Some(false)` still counts as "the check ran", while this one + // counts only widths where the check ran *and* the pool passed it + // with room to spare. A single flag would let a sweep in which + // every verdict was negative report the same coverage as one in + // which every verdict was positive. assert!( saw_realized_placement_check || !crate::decode_affinity::pinning_supported() @@ -21086,6 +21269,426 @@ mod tests { } } + /// The widest swept width that can still be one-per-core on this host, or + /// `None` when the host cannot answer the question at all. + /// + /// Returning `None` for an unsupported platform is a *skip with a reason*; + /// it is never conflated with a passing answer, and every caller says out + /// loud which of the two it got. + fn placement_probe_width() -> Option<(usize, usize)> { + if !crate::core_topology::DETECTION_SUPPORTED + || !crate::decode_affinity::pinning_supported() + || !crate::decode_affinity::affinity_observation_supported() + { + return None; + } + let allowed = crate::decode_affinity::allowed_cpus()?; + let topology = crate::core_topology::require_host_for_placement() + .expect("DETECTION_SUPPORTED targets must resolve a topology or panic"); + let cores = topology.physical_cores_within(&allowed); + // Leave the dispatcher its reserved CPU, so the requested width is one + // the pool can actually realize one-per-core. + let width = cores.saturating_sub(1).min(8); + (width >= 2).then_some((width, cores)) + } + + /// A misspelled placement must fall back to the default, say so, and still + /// decode. + /// + /// End-to-end rather than a unit test of the parser, because the two ways + /// this knob can be inert are both outside the parser: reading the wrong + /// variable name, and swallowing the diagnostic. Either one produces a knob + /// a user can set with no error and no effect, which is #1792's shape and + /// the failure `verify_documented_env_vars.py` exists to catch statically. + #[test] + #[cfg_attr(miri, ignore = "spawns a child process")] + fn a_misspelled_placement_falls_back_loudly_and_still_builds_a_pool() { + let width = "2"; + let output = std::process::Command::new(std::env::current_exe().unwrap()) + .arg("--exact") + .arg("kernels::matmul_nbits::tests::spmd_realized_width_subprocess") + .arg("--nocapture") + .arg("--test-threads=1") + .env(SPMD_WIDTH_CHILD_ENV, width) + .env(DECODE_THREADS_ENV, width) + .env("RAYON_NUM_THREADS", width) + .env(crate::decode_spmd::PERSISTENT_POOL_ENV, "1") + .env(crate::decode_affinity::DECODE_PLACEMENT_ENV, "one-per-core") + .env_remove(crate::decode_affinity::DECODE_AFFINITY_ENV) + .env_remove(crate::decode_spmd::DECODE_SCHEDULE_ENV) + .env_remove(SPMD_PARITY_CHILD_ENV) + .env_remove(PLACEMENT_DISHONEST_ENV) + .output() + .expect("run the misspelled-placement child process"); + let stdout = String::from_utf8_lossy(&output.stdout); + let stderr = String::from_utf8_lossy(&output.stderr); + + // Not fatal: this knob ranks a set already chosen, so refusing to + // decode over a misspelled tuning value would be the worse failure. + assert!( + output.status.success(), + "a misspelled placement aborted the child instead of falling back \ + (status={}):\nstdout:\n{stdout}\nstderr:\n{stderr}", + child_status_detail(&output.status) + ); + assert!( + stdout.contains(&format!( + "policy={}", + crate::decode_affinity::CorePlacement::default().as_str() + )), + "the child did not fall back to the default placement:\n{stdout}" + ); + // ...and never silent. A fallback nobody is told about is how a user + // ends up believing a knob is in force for the life of a deployment. + assert!( + stderr.contains(crate::decode_affinity::DECODE_PLACEMENT_ENV) + && stderr.contains("one-per-core"), + "the fallback was silent -- stderr must name the variable and the value it \ + rejected:\nstderr:\n{stderr}" + ); + } + + /// The code returned when no layout can be predicted at all. + /// + /// Deliberately not one of the real codes: a degenerate input must not + /// return the value that happens to make a caller's assertion pass, which + /// is the false-oracle shape this change exists to remove. + const UNPREDICTABLE_PLACEMENT: &str = "unpredictable"; + + /// The placement the named policy predicts for a `workers`-wide pool on + /// `allowed`, expressed in the same wire codes the child reports. + /// + /// Exact rather than approximate, and that matters: `node_shards_with` + /// orders the *whole* allowed set through `order_pin_targets_for` and then + /// caps only the worker *count* (`reserve_single_group_headroom` returns a + /// count, it does not remove CPUs from the list), and the pool pins worker + /// `i` to `cpus[i % len]`. So the first `workers` entries of that ordering + /// are the CPUs the workers actually get. + /// + /// Comparing this against `realized` is not circular. The prediction comes + /// from the ordering policy; `realized` comes from each worker reading the + /// affinity mask the kernel is actually enforcing for it. The claim is that + /// those two agree. + fn predicted_placement_code( + allowed: &[usize], + topology: &crate::core_topology::CoreTopology, + policy: crate::decode_affinity::CorePlacement, + workers: usize, + ) -> &'static str { + let ordered = + crate::decode_affinity::order_pin_targets_for(allowed, Some(topology), policy); + if ordered.is_empty() || workers == 0 { + return UNPREDICTABLE_PLACEMENT; + } + let taken: Vec = (0..workers).map(|i| ordered[i % ordered.len()]).collect(); + if topology.physical_cores_within(&taken) < workers { + "shared-core" + } else { + "one-per-core" + } + } + + /// The cpuset the child measured, together with the topology, when the + /// parent can still see the same set the child did. + /// + /// The child inherits this process's affinity mask, so the two agree by + /// construction. Checking rather than assuming is the point: if they ever + /// diverge, every layout this parent predicts is about a different machine + /// than the one the child measured. The comparison is on the *set*, not on + /// its cardinality -- equal-sized but differently populated cpusets predict + /// different layouts, so a size check would not be the guard this claims to + /// be. + /// + /// A divergence is a skip with a printed reason rather than a failure, and + /// that is deliberate. A cgroup resize or a CPU hot-unplug between the + /// child's read and the parent's is an environmental event, not a defect in + /// the code under test, and a red whose cheapest resolution is to delete the + /// check is worse than a stated skip. + fn parent_cpuset_matching( + report: &RealizedWidth, + ) -> Option<(Vec, &'static crate::core_topology::CoreTopology)> { + let topology = crate::core_topology::require_host_for_placement() + .expect("DETECTION_SUPPORTED targets must resolve a topology or panic"); + let allowed = crate::decode_affinity::allowed_cpus().unwrap_or_default(); + if allowed.is_empty() || allowed != report.allowed_cpus { + eprintln!( + "skipping the layout prediction: the child measured cpuset {:?} and this \ + parent now sees {allowed:?}, so the machine changed shape under the test \ + and any layout predicted here would describe a different one ({report:?})", + report.allowed_cpus + ); + return None; + } + Some((allowed, topology)) + } + + /// With the spread policy **selected**, the workers must land one per + /// physical core. + /// + /// This is the assertion the default sweep used to make unconditionally. + /// It is correct here and wrong there for one reason: here the policy is + /// named, so the test asserts a claim it actually asked for. #1802 may + /// change what the default is; it cannot change what `spread` means. + #[test] + #[cfg_attr(miri, ignore = "spawns a child process")] + fn an_explicit_spread_policy_places_one_worker_per_physical_core() { + let Some((probe, cores)) = placement_probe_width() else { + eprintln!( + "skipping explicit-spread placement check: this target cannot pin, cannot \ + read a mask back, or has fewer than 3 allowed physical cores" + ); + return; + }; + // Every published width the host has cores for, not one probe width. + // The assertions this replaces ran at every swept width, and narrowing + // them to a single one would leave a pool-side pinning defect that only + // appears past 8 workers unguarded -- placement wrong, honesty green, + // nothing red. + let mut widths: Vec = BENCHMARKED_DECODE_WIDTHS + .iter() + .copied() + .filter(|width| (2..=cores).contains(width)) + .collect(); + if !widths.contains(&probe) { + widths.push(probe); + } + let mut checked = 0usize; + for width in widths { + let report = realized_width_report_with( + width, + false, + Some(crate::decode_affinity::CorePlacement::Spread), + ); + assert_eq!( + report.placement_policy, "spread", + "the child did not receive the policy this test selected, so it is \ + measuring something else ({report:?})" + ); + if report.nodes > 1 { + eprintln!( + "width {width}: skipping the explicit-spread check because a \ + NUMA-split pool rebalances workers per node, so the single-group \ + layout this asserts is not the one that ran ({report:?})" + ); + continue; + } + assert!( + report.pool_built && report.workers >= 2 && report.workers <= cores, + "width {width} did not produce a pool inside the core budget, so the \ + assertion below would be about the host rather than the policy \ + ({report:?})" + ); + assert_eq!( + report.realized, "one-per-core", + "`spread` was selected explicitly and the host has room for it at width \ + {width}, yet the workers observed `{}` for themselves -- the selector is \ + inert ({report:?})", + report.realized + ); + assert_eq!( + report.planned_distinct_cores, + Some(true), + "the planner did not even intend one worker per physical core under an \ + explicit `spread` at width {width} ({report:?})" + ); + // Honest under the policy too: a spread that is reported but not + // applied is the #1792 defect wearing a policy name. + assert_ne!( + report.placement_honest, + Some(false), + "width {width}: the pool reports CPUs the kernel did not enforce \ + ({report:?})" + ); + checked += 1; + } + assert!( + checked > 0, + "every width skipped, so `spread` was never actually asserted -- the probe \ + width said this host could realize one worker per core and then no width did" + ); + } + + /// With the compact policy selected, workers share cores -- and the sweep's + /// policy-neutral assertions must still all pass. + /// + /// Two claims in one child, deliberately. The first is that the selector + /// does something (a knob that parses and changes nothing is #1792 again). + /// The second is the mutation #1802 needs: everything the default sweep + /// asserts has to hold when the placement is compact, or the sweep is still + /// encoding the spread policy somewhere. + #[test] + #[cfg_attr(miri, ignore = "spawns a child process")] + fn a_compact_policy_shares_cores_and_still_passes_every_policy_neutral_check() { + let Some((width, _cores)) = placement_probe_width() else { + eprintln!( + "skipping compact placement check: this target cannot pin, cannot read a \ + mask back, or has fewer than 3 allowed physical cores" + ); + return; + }; + let report = realized_width_report_with( + width, + false, + Some(crate::decode_affinity::CorePlacement::Compact), + ); + assert_eq!( + report.placement_policy, "compact", + "the child did not receive the policy this test selected ({report:?})" + ); + + // Every policy-neutral claim the default sweep makes, restated here + // against a compact pool. If any of these fail, the sweep has not been + // made policy-neutral and #1802 is still blocked. + assert_ne!( + report.placement_honest, + Some(false), + "a compact pool reports CPUs the kernel did not enforce ({report:?})" + ); + assert_ne!( + report.realized, "affinity-query-failed", + "a worker could not read its own mask ({report:?})" + ); + assert!( + !crate::decode_affinity::affinity_observation_supported() + || report.realized != "no-affinity-query", + "the realized-placement observation is switched off ({report:?})" + ); + assert!( + report.fully_pinned, + "a compact pool left workers unpinned, so `compact` is being read as \ + `off` rather than as a layout ({report:?})" + ); + + // ...and the claim that makes the selector worth having: the workers + // are where the compact ordering says they are. + // + // The predicate is derived from *this cpuset*, not from the machine. + // An earlier version demanded `shared-core` whenever the host had SMT + // anywhere, which is a different question: `taskset -c 0,2,4,6` on an + // SMT host leaves no sibling pair inside the allowed set, so compact + // and spread are the same layout there and the demand would have failed + // a legitimate host. Predicting from the allowed set answers what the + // policy will actually do here, so the assertion holds on every cpuset + // and still catches an inert selector on the ones that can show it. + if report.nodes > 1 { + eprintln!( + "skipping the compact layout prediction: a NUMA-split pool rebalances \ + workers per node, so the single-group ordering predicted here is not the \ + one that ran ({report:?})" + ); + return; + } + assert!( + report.workers >= 2, + "a one-worker pool cannot share a core with anything, so the layout claim \ + below would be about the host rather than the policy ({report:?})" + ); + let Some((allowed, topology)) = parent_cpuset_matching(&report) else { + return; + }; + let compact_predicted = predicted_placement_code( + &allowed, + topology, + crate::decode_affinity::CorePlacement::Compact, + report.workers, + ); + assert_ne!( + compact_predicted, UNPREDICTABLE_PLACEMENT, + "no layout could be predicted for a pool this parent has already confirmed \ + has workers and a cpuset ({report:?})" + ); + let spread_predicted = predicted_placement_code( + &allowed, + topology, + crate::decode_affinity::CorePlacement::Spread, + report.workers, + ); + assert_eq!( + report.realized, compact_predicted, + "the compact ordering puts {} workers on {compact_predicted}, but the workers \ + observed `{}` for themselves -- either the selector parses and changes \ + nothing (the #1792 shape) or the pin did not take ({report:?})", + report.workers, report.realized + ); + if compact_predicted == spread_predicted { + eprintln!( + "the discrimination half of the compact check did not run: on this \ + {}-CPU / {}-core cpuset both policies predict `{compact_predicted}` at \ + {} workers, so no layout difference exists to observe", + report.allowed, report.cores, report.workers + ); + } else { + assert_ne!( + report.realized, spread_predicted, + "compact was selected but the pool landed in the layout `spread` would \ + have produced, so the selector is inert ({report:?})" + ); + } + } + + /// The mutation for the explicit-spread arm: a compact placement must + /// *fail* the spread assertion. + /// + /// Without this, `an_explicit_spread_policy_places_one_worker_per_physical_core` + /// could be passing because `realized` is hardcoded to `one-per-core` + /// somewhere, and nobody would know. It re-runs the same child under the + /// other policy and checks the same predicate answers differently -- so the + /// predicate is shown to discriminate, on this host, today. + #[test] + #[cfg_attr(miri, ignore = "spawns a child process")] + fn the_spread_assertion_rejects_a_compact_placement() { + let Some((width, _cores)) = placement_probe_width() else { + eprintln!("skipping spread-assertion mutation: unsupported target"); + return; + }; + let compact = realized_width_report_with( + width, + false, + Some(crate::decode_affinity::CorePlacement::Compact), + ); + if compact.nodes > 1 { + eprintln!( + "skipping spread-assertion mutation: a NUMA-split pool rebalances workers \ + per node, so the single-group ordering predicted here is not the one that \ + ran ({compact:?})" + ); + return; + } + let Some((allowed, topology)) = parent_cpuset_matching(&compact) else { + return; + }; + // Whether a mutation exists at all is a property of *this cpuset*, not + // of the machine: `taskset -c 0,2,4,6` on an SMT host admits no sibling + // pair, so compact and spread are the same layout and there is nothing + // to mutate. Keying this off whole-machine `has_smt()` would have + // demanded a difference the cpuset cannot produce. + // + // Written as "only proceed on a positive shared-core prediction" rather + // than "skip on one-per-core", so the unpredictable case skips too + // instead of falling into the assertion. + if predicted_placement_code( + &allowed, + topology, + crate::decode_affinity::CorePlacement::Compact, + compact.workers, + ) != "shared-core" + { + eprintln!( + "skipping spread-assertion mutation: this {}-CPU / {}-core cpuset has no \ + SMT sibling pair among the first {} compact slots, so `compact` is the \ + same layout as `spread` here and no mutation exists to apply", + compact.allowed, compact.cores, compact.workers + ); + return; + } + assert_ne!( + compact.realized, "one-per-core", + "a compact placement satisfied the one-per-core predicate, so that \ + predicate cannot fail and the explicit-spread test proves nothing \ + ({compact:?})" + ); + } + #[test] fn spmd_adaptive_calibrated_decode_is_bit_identical_to_flat() { // Token-exactness of the adaptive path: with diff --git a/docs/performance/numa-decode-plan.md b/docs/performance/numa-decode-plan.md index cb7d2df267..15fdc7881a 100644 --- a/docs/performance/numa-decode-plan.md +++ b/docs/performance/numa-decode-plan.md @@ -69,6 +69,34 @@ Qwen2.5-Coder-7B int4, 32 decode threads, steady M=1, 5 runs x 3 rounds): jitter that makes the unpinned pool swing run to run. Greedy token ids are bit-identical with and without pinning (it only changes placement, not math). +### Placement inside the chosen CPU set: `ONNX_GENAI_CPU_DECODE_PLACEMENT` + +`ONNX_GENAI_CPU_DECODE_AFFINITY` selects *which* CPUs the decode pool may use. +`ONNX_GENAI_CPU_DECODE_PLACEMENT` selects how the workers are arranged inside +that set. The two are orthogonal and can be combined. + +- unset / `spread` -- **default**. Worker `i` takes a distinct physical core for + as long as cores last, only then doubling up on SMT siblings. Fastest on a + host the process has to itself. +- `compact` -- SMT siblings are filled before the next core is used, so an + N-worker pool occupies about `N / siblings` physical cores and leaves the rest + of the machine to a co-tenant. + +Neither is universally better, which is why it is a selector and not a constant. +Measured with four DRAM-bandwidth hogs on a 16-core/32-thread host, `compact` +ran 4.54 ms/token against `spread`'s 5.03--6.26; with eight hogs covering both +halves the ranking inverts (`compact` 15.92 against `spread` 6.08). The full +table is in `crates/onnx-runtime-ep-cpu/src/core_topology.rs`'s module docs. + +An unrecognized value is reported on stderr once and then treated as `spread`, +rather than refusing to decode: unlike `ONNX_GENAI_CPU_DECODE_AFFINITY`, which +decides which CPUs the process may touch, this knob only ranks a set already +chosen, so a misspelling must not be fatal. It is never silent. + +Because the ordering is only a ranking of a permitted set, both values are +permutations of the same CPU list: no placement policy can widen or shrink the +CPUs a cgroup or `taskset` allows. + Not yet implemented (steps 4--5, deferred): sharding a single projection across node-local replicas on both sockets. A naive dual-node pool regresses badly -- a 64-thread pool spanning both sockets with interleaved memory measured