From 8fc8513f608ce7b795c75e9f64b39d664dcbf192 Mon Sep 17 00:00:00 2001 From: Nils Homer Date: Sat, 12 Sep 2026 18:23:54 -0700 Subject: [PATCH 1/2] feat(pipeline-core): park idle pool workers to fix oversubscription cliff MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The chain-builder pool is replicated round-robin polling with no producer->consumer wakeup: an idle worker (whole dispatch pass did no work) exponential-backoff *sleeps*, then re-polls. On an oversubscribed pipeline (more workers than available parallel work) the surplus workers churn poll->sleep->poll, and that churn actively slows the productive workers. On `fgumi correct` at 32 threads this was a hard cliff: ~6.5s wall vs ~1.9s at 16 threads (pool 92% idle, 208s of aggregate idle churn), because a single serial grouping stage feeds the parallel work and nothing else can run. Two layers: L0 — take the cheap wins, remove the contention: - Test-and-test-and-set on the Serial step mutex. `parking_lot::Mutex::try_lock` is a compare_exchange that writes (invalidates) the holder's cache line even when it fails; a surplus worker probing a held `Serial + Affinity::None` step (the throughput ceiling) stole that line every idle pass. Probe `is_locked()` (a shared load) first and report `Contention` without the write. - Raise the pool backoff floor 1us -> 20us so a surplus worker's stray-success reset does not drop it back onto a per-microsecond re-poll ramp. Wake latency for a worker that made progress is unaffected (it does not sleep). L2 — park the surplus (the real fix): - `PoolEventCount`, a Vyukov event-count over parking_lot Mutex/Condvar. An idle worker arms (`prepare_wait`: register + SeqCst fence + snapshot generation), re-polls its real condition one more pass, then blocks (`wait`) only if still empty; a producer publishes work then `notify_one` (SeqCst fence + a read-shared `waiters` load; when nobody is parked that is the whole cost — no store, no lock, no wake). Correctness is the C++20 SeqCst-fence rule: the arm fence and the notify fence are totally ordered, so no interleaving parks on a poppable item. Notify-under-lock gives the condvar the token a bare condvar lacks. - Notify seam is the driver's dispatch outcome, not the queue: `Progress` (pushed or held an item, including a reorder must-accept stash) -> notify_one; `Finished` (an output edge closed) -> notify_all. This captures the ordinal-unblocking stash without threading the parker through every queue. - Concurrency ceiling: notify_one wakes a parked worker only while `awake = n_pool - waiters < ceiling`, so the per-Progress cascade cannot re-wake the whole surplus. Default ceiling = n_pool (ungated); a controller can lower it later. notify_all (cancel/close/Finished) is never gated. - Pinned workers (Exclusive/sticky owners, and Serial affinity targets) keep the old sleep_backoff: a shared notify_one cannot target a specific worker, so a push meant for a pinned step's owner could wake a peer that Skips it. They are a handful of N and were never the problem. - Detached driver threads participate as notifiers (their pushes wake pool waiters via the same dispatch seam) but not as waiters (they keep their Park backoff). - Cancel/error wake every parked worker: `PipelineSignal` holds a late-bound Weak (bound in `run` before any spawn, no Arc cycle) and calls notify_all after the terminal-state CAS in both `cancel` and `record_error` (a step Err calls only record_error, never cancel), so a parked worker observes is_done() instead of sleeping to its timeout. The wait deadline (the existing backoff) remains as a self-heal. - Only wired when `n_threads > 1`; the single-thread and fused paths pass `None` and keep the existing sleep-backoff idle (behaviour-identical off path). Result on `fgumi correct` at 32 threads: ~1.9s wall (from ~6.5s), flat vs 16 threads; aggregate idle 208s -> 35s (real parking, not spin-churn). `extract` K=2 (push-heavy) and `sort` (exercises the detached-driver notifier + pinned paths) show no regression. --- crates/fgumi-pipeline-core/src/builder.rs | 50 ++ .../src/runtime/detached.rs | 12 + .../fgumi-pipeline-core/src/runtime/driver.rs | 219 ++++++++ .../src/runtime/event_count.rs | 491 ++++++++++++++++++ crates/fgumi-pipeline-core/src/runtime/mod.rs | 2 + .../src/runtime/worker_core.rs | 21 +- crates/fgumi-pipeline-core/src/signal.rs | 40 ++ 7 files changed, 834 insertions(+), 1 deletion(-) create mode 100644 crates/fgumi-pipeline-core/src/runtime/event_count.rs diff --git a/crates/fgumi-pipeline-core/src/builder.rs b/crates/fgumi-pipeline-core/src/builder.rs index 44df0c3e5..d7b8cb455 100644 --- a/crates/fgumi-pipeline-core/src/builder.rs +++ b/crates/fgumi-pipeline-core/src/builder.rs @@ -1222,6 +1222,32 @@ impl Pipeline { // with Serial+sticky+Affinity targeting). Indexed by worker id. let sticky_owners = assign_sticky_owners(&steps, &owners, n_threads); + // 1b. Compute which workers are "pinned" — the sole eligible dispatcher + // of some step: an Exclusive owner, a sticky owner, or the affinity + // target of a `Serial` step (including a non-sticky affinity like the + // `correct` chain's pinned `GroupByQueryname`). A pinned worker must NOT + // deep-park on the shared event-count: `notify_one` cannot target a + // specific worker, so a push meant to wake a pinned step's owner could + // wake a peer that `Skip`s it, leaving the owner asleep. Pinned workers + // keep the exponential-backoff sleep (they are ≤ a handful of N and were + // never the oversubscription problem). Computed while `steps` is still + // alive, before `build_worker_storage` consumes it. + let pinned_workers: Vec = { + let mut pinned = vec![false; n_threads]; + for (w, p) in pinned.iter_mut().enumerate() { + *p = owners.contains(&Some(w)) || sticky_owners[w].is_some(); + } + for step in &steps { + if step.kind() == crate::step::StepKind::Serial + && let Some(target) = step.affinity().target_worker(n_threads) + && target < n_threads + { + pinned[target] = true; + } + } + pinned + }; + // 2. Build per-step contexts (input + output handles). Domain counters // are allocated only when telemetry is on (`counters_enabled`), so the // telemetry-off path keeps its byte-identical zero-allocation route — @@ -1309,6 +1335,21 @@ impl Pipeline { let signal_arc = Arc::clone(&signal); + // 4-0. Pool event-count parker (Layer 2 oversubscription fix). Only for + // the multi-worker path: with one worker there is no peer to notify and + // nothing to park behind. Bound to the signal *before* any worker spawns + // so a terminal transition (cancel/error) can wake every parked worker. + // The signal holds only a `Weak`, so this `Arc` (kept alive here for the + // whole run scope) is what keeps the parker live while workers could be + // parked. + let parker: Option> = if n_threads > 1 { + let ec = Arc::new(crate::runtime::event_count::PoolEventCount::new(n_threads)); + signal_arc.bind_event_count(&ec); + Some(ec) + } else { + None + }; + // 4a. Optional deadlock-detection monitor. Spawns a watcher // thread that periodically samples the stats snapshot; if no // step's `progress + finished` counter has advanced for @@ -1542,6 +1583,7 @@ impl Pipeline { let liveness_clone = Arc::clone(&liveness); // Detached drivers take board slots after the pool workers. let board_clone = worker_board.clone(); + let parker_clone = parker.clone(); let state_slot = n_threads + driver_idx; let thread_name = match group.label() { DetachedGroup::Shared(label) => format!("fgumi-driver-{label}"), @@ -1568,6 +1610,7 @@ impl Pipeline { &liveness_clone, board_clone.as_deref(), state_slot, + parker_clone.as_deref(), ); })) { @@ -1630,6 +1673,9 @@ impl Pipeline { // the zero-cost path. worker_board.as_deref(), 0, + // n_threads == 1: no parker (nobody to notify / park behind). + None, + pinned_workers[0], ); })) { signal_arc.cancel(); @@ -1654,6 +1700,8 @@ impl Pipeline { // Pool worker `worker_id` takes board slot `worker_id` (telemetry // on); `None` on the zero-cost path. let board_clone = worker_board.clone(); + let parker_clone = parker.clone(); + let pinned = pinned_workers[worker_id]; let handle = thread::Builder::new() .name(format!("fgumi-worker-{worker_id}")) @@ -1683,6 +1731,8 @@ impl Pipeline { scheduler_clone.as_ref(), board_clone.as_deref(), worker_id, + parker_clone.as_deref(), + pinned, ); })) { diff --git a/crates/fgumi-pipeline-core/src/runtime/detached.rs b/crates/fgumi-pipeline-core/src/runtime/detached.rs index e2426217c..ad150aa22 100644 --- a/crates/fgumi-pipeline-core/src/runtime/detached.rs +++ b/crates/fgumi-pipeline-core/src/runtime/detached.rs @@ -337,6 +337,12 @@ pub fn run_detached_driver( liveness: &crate::liveness::LivenessCounter, board: Option<&crate::runtime::worker_state::WorkerStateBoard>, state_slot: usize, + // The pool's event-count, so this driver's productive pushes wake parked + // pool workers (the notify seam in `dispatch_one_step`). Passed with + // `pinned = true` below so the driver itself keeps its `Park` backoff and + // never deep-parks on the shared event-count — it is a notifier, not a + // waiter (spec §5). `None` when the pool has no parker (`n_threads == 1`). + parker: Option<&crate::runtime::event_count::PoolEventCount>, ) { let primary = group.primary_step(); let mut row = build_driver_storage(group.into_steps(), contexts.inputs.len()); @@ -352,6 +358,11 @@ pub fn run_detached_driver( &DrainFirstScheduler, board, state_slot, + parker, + // `pinned`: a detached driver is not a pool waiter. This routes its idle + // to the `Park` backoff (its existing behaviour) rather than the shared + // event-count wait, while still passing `parker` above for the notify. + true, ); } @@ -496,6 +507,7 @@ mod tests { &crate::liveness::LivenessCounter::new(1), None, 0, + None, ); } diff --git a/crates/fgumi-pipeline-core/src/runtime/driver.rs b/crates/fgumi-pipeline-core/src/runtime/driver.rs index 6c62e006e..80fe91c8c 100644 --- a/crates/fgumi-pipeline-core/src/runtime/driver.rs +++ b/crates/fgumi-pipeline-core/src/runtime/driver.rs @@ -30,6 +30,7 @@ use crate::erased::ErasedStepCtx; use crate::liveness::LivenessCounter; use crate::runtime::contexts::ChainContexts; use crate::runtime::drain::StepDrainCounter; +use crate::runtime::event_count::PoolEventCount; use crate::runtime::live::LiveSteps; use crate::runtime::scheduler::{Scheduler, WalkDirection}; use crate::runtime::stats::PipelineStats; @@ -68,6 +69,11 @@ const STICKY_BURST_LIMIT: usize = 1024; /// `board` — optional per-thread state board for scheduling telemetry; `None` /// keeps the loop's stamping a no-op (telemetry-off parity). /// `state_slot` — this thread's slot in `board` (ignored when `board` is `None`). +/// `parker` — the pool's event-count. `Some` only for pool workers when +/// `n_threads > 1`; drives both the idle-branch park (instead of `sleep_backoff`) +/// and the on-`Progress`/`Finished` notify that wakes parked peers. `None` +/// (single-thread paths, and pinned workers — see the idle branch) keeps the +/// existing `sleep_backoff` behaviour. #[allow(clippy::too_many_arguments)] // per-step shared state plus the liveness shard; a struct would only rename it #[allow(clippy::too_many_lines)] // the whole-pass sticky + round-robin + backoff discipline is documented inline; splitting it would scatter the invariant across functions pub fn run_worker_loop( @@ -81,7 +87,20 @@ pub fn run_worker_loop( scheduler: &dyn Scheduler, board: Option<&WorkerStateBoard>, state_slot: usize, + parker: Option<&PoolEventCount>, + pinned: bool, ) { + // A pinned worker (the sole eligible dispatcher of an Exclusive/sticky step, + // or the affinity target of a Serial step) must NOT deep-park on the shared + // event-count: `notify_one` cannot target a specific worker, so a push meant + // to wake *this* worker's pinned step could wake a peer that `Skip`s it, + // leaving the pinned step's owner asleep to its timeout. Pinned workers keep + // `sleep_backoff`. Its step is also often driven by out-of-engine events + // (a reader's prefetch thread, the affinity-pinned grouper) that no push + // notifies anyway. `parker` is threaded to `None` for the idle branch here, + // while still being used for the on-Progress notify (a pinned worker's + // pushes must still wake parked peers). + let idle_parker = if pinned { None } else { parker }; // Per-worker worklist of still-dispatchable steps, in chain order. A step // is removed when it returns `StepOutcome::Finished`; the worker exits once // the list is empty. Build-time `Skip` placeholders (Exclusive steps owned @@ -156,6 +175,7 @@ pub fn run_worker_loop( is_driver, board, state_slot, + parker, ) else { break; // Skip }; @@ -203,6 +223,7 @@ pub fn run_worker_loop( is_driver, board, state_slot, + parker, ); did_work |= outcome.did_work; if outcome.removed_sticky_owner { @@ -233,7 +254,85 @@ pub fn run_worker_loop( worker.reset_backoff(); } else if signal.is_done() { break; + } else if let Some(ec) = idle_parker { + // Event-count park (Layer 2): block instead of spin-poll when the + // whole pass found no work. The two-phase protocol closes the + // lost-wakeup race: arm (register + fence) → re-poll the real + // condition → block only if still empty. + // + // Arm and re-poll *before* stamping `Parked` or starting the idle + // timer: a productive re-poll (work appeared in the arm→wait window) + // must not be recorded as idle, nor leave the board reading `Parked` + // while the re-poll's `dispatch_one_step` runs real work. + + // Phase 1: arm. The fence in `prepare_wait` orders the waiter + // registration before the re-poll below, so a producer that + // publishes work after this point either sees us as a waiter (and + // bumps the generation) or we see its item on the re-poll. + let key = ec.prepare_wait(); + + // Phase 2: re-poll the *real* condition — one full dispatch pass. + // This is one pass per park episode (amortised over the block), not + // per iteration, so it is not the continuous polling Layer 2 removes. + let recheck = round_robin_dispatch( + entries, + &mut live, + worker.sticky_owner, + contexts, + drain_counters, + signal, + stats, + liveness, + worker.thread_id, + scheduler.walk(), + is_driver, + board, + state_slot, + parker, + ); + if recheck.removed_sticky_owner { + sticky_live = false; + } + + if recheck.did_work || signal.is_done() { + // Work appeared in the arm→wait window (or we are shutting + // down): do not block, and record no idle — the re-poll was + // productive. Balance the arm and re-ramp fresh. + ec.cancel_wait(key); + worker.reset_backoff(); + } else { + // No progress remains: only now do we actually park. Stamp + // `Parked` and start the idle timer immediately before blocking, + // so the idle attribution covers exactly the wait, not the + // preceding (possibly productive) re-poll. + if let Some(b) = board { + b.stamp(state_slot, crate::runtime::worker_state::WorkerState::Parked, None); + } + let sleep_start = stats.map(|_| Instant::now()); + + // Phase 3: block until the generation moves (a peer's notify), + // the deadline elapses (self-heal), or a spurious wake. Keep the + // exponential backoff as the wait *timeout* so a missed wakeup + // cannot hang longer than the cap, and the deadlock monitor + // still ticks. Re-poll on return regardless of outcome. + let deadline = worker.backoff_deadline(); + let _ = ec.wait(key, deadline); + worker.increase_backoff(); + + if let (Some(stats), Some(ss)) = (stats, sleep_start) { + let ns = u64::try_from(ss.elapsed().as_nanos()).unwrap_or(u64::MAX); + match worker.role() { + WorkerRole::Pool => stats.record_worker_idle(worker.thread_id, ns), + WorkerRole::Driver { primary_step } => { + stats.record_detached_idle(primary_step, ns); + stats.record_detached_park(primary_step); + } + } + } + } } else { + // No parker (single-thread path or a pinned worker): keep the + // original exponential-backoff sleep. if let Some(b) = board { b.stamp(state_slot, crate::runtime::worker_state::WorkerState::Parked, None); } @@ -297,6 +396,7 @@ fn round_robin_dispatch( is_driver: bool, board: Option<&WorkerStateBoard>, state_slot: usize, + parker: Option<&PoolEventCount>, ) -> RoundRobinOutcome { let mut did_work = false; // Steps that finished this pass, removed from `live` after the walk. A @@ -334,6 +434,7 @@ fn round_robin_dispatch( is_driver, board, state_slot, + parker, ) else { continue; // Skip (build-time placeholder; should not appear in `live`) }; @@ -431,6 +532,14 @@ fn dispatch_one_step( // slot in `board` (ignored when `board` is `None`). board: Option<&WorkerStateBoard>, state_slot: usize, + // The pool's event-count. On a `Progress` outcome (the step pushed or held + // an item — including a reorder must-accept stash) wake one parked peer; on + // `Finished` (an output edge just closed) wake all. `None` for single-thread + // paths. This is the notify seam: `Progress`/`Finished` are exactly the + // outcomes that can give a parked consumer new work, and using the outcome + // (rather than the `BranchOutputHandle::push` site) captures the stash-then- + // return-`Ok` reorder case without threading the parker through every queue. + parker: Option<&PoolEventCount>, ) -> Option { let outputs_any: &(dyn Any + Send + Sync) = contexts.outputs[step_idx.0].as_ref(); let mut ctx = ErasedStepCtx { @@ -492,6 +601,22 @@ fn dispatch_one_step( result: Ok(StepOutcome::Finished), name: "", }) + } else if step.is_locked() { + // Test-and-test-and-set: probe with a shared-mode load before + // attempting the lock. `parking_lot::Mutex::try_lock` is a + // `compare_exchange` on the lock word that writes (and so + // invalidates) the holder's cache line *even when it fails*. A + // surplus worker that `try_lock`s a `Serial + Affinity::None` + // step held by the productive worker steals that line on every + // idle pass, and that step is often the throughput ceiling + // (e.g. `GroupByQueryname` in `correct`). `is_locked()` is a + // plain relaxed load — no line-ownership transfer — so probing + // it first turns the common "held by a peer" case into a + // read-only `Contention` with no write to the holder's line. + Some(DispatchInfo { + result: Ok(StepOutcome::Contention), + name: "", + }) } else { match step.try_lock() { None => Some(DispatchInfo { @@ -528,6 +653,20 @@ fn dispatch_one_step( && matches!(i.result, Ok(StepOutcome::Progress | StepOutcome::Finished)) { liveness.bump(worker_slot); + // Notify seam: a productive dispatch may have given a parked consumer + // new work. `Progress` = "pushed or held an item" (incl. a reorder + // must-accept stash that unblocks an ordinal waiter) → wake one parked + // peer, which cascades (its own downstream push wakes the next). + // `Finished` closed an output edge (a parked consumer must observe + // end-of-stream, and a sink's `Finished` closes no edge but its peers + // may be waiting on the drain latch) → wake all. When nobody is parked, + // both are a single relaxed load (see `PoolEventCount`). + if let Some(ec) = parker { + match i.result { + Ok(StepOutcome::Finished) => ec.notify_all(), + _ => ec.notify_one(), + } + } } if let (Some(stats), Some(start)) = (stats, start) { @@ -595,6 +734,8 @@ mod tests { &crate::runtime::scheduler::ChainOrderScheduler, None, 0, + None, + false, ); // If we reach this line, the loop exited cleanly. } @@ -635,6 +776,8 @@ mod tests { &crate::runtime::scheduler::ChainOrderScheduler, Some(&board), 0, + None, + false, ); assert_eq!(board.read(0).0, WorkerState::Parked); } @@ -843,6 +986,7 @@ mod tests { false, None, 0, + None, ) .unwrap(); assert!(matches!(info.result, Ok(StepOutcome::Finished))); @@ -887,6 +1031,7 @@ mod tests { false, None, 0, + None, ) .unwrap(); assert!(matches!(info1.result, Ok(StepOutcome::Finished))); @@ -908,6 +1053,7 @@ mod tests { false, None, 0, + None, ) .unwrap(); assert!(matches!(info2.result, Ok(StepOutcome::Finished))); @@ -922,6 +1068,67 @@ mod tests { ); } + /// Layer 0 (test-and-test-and-set): when a `Serial` step's mutex is already + /// held (by the productive worker), a surplus worker's dispatch reports + /// `Contention` via the `is_locked()` probe and does NOT run the step. The + /// probe is a shared-mode load, so it never takes the lock word's write + /// path (`try_lock`'s `compare_exchange`) that would invalidate the + /// holder's cache line. We prove the observable contract: the step body is + /// not entered while the lock is held. + #[test] + fn serial_contention_probe_does_not_run_held_step() { + let (steps, graph, runs) = three_step_chain(StepKind::Serial); + let contexts = Arc::new(build_chain_contexts( + &steps, + &graph, + crate::builder::InstrumentationLevel::Off, + false, + )); + let mid = StepIdx(1); + let counter = StepDrainCounter::new(1); + let signal = PipelineSignal::new(); + + let shared = Arc::new(Mutex::new(steps.into_iter().nth(1).unwrap())); + let drain = Arc::new(DrainGate::default()); + + // Simulate the productive worker holding the step's mutex. + let held = Arc::clone(&shared); + let guard = held.lock(); + assert!(shared.is_locked(), "precondition: the Serial mutex is held"); + + // A surplus worker dispatches the same step while it is held. + let mut entry = + WorkerStepEntry::Shared { step: Arc::clone(&shared), drain: Arc::clone(&drain) }; + let info = dispatch_one_step( + &mut entry, + mid, + &contexts, + &counter, + &signal, + None, + &LivenessCounter::new(1), + 0, + false, + None, + 0, + None, + ) + .unwrap(); + + assert!( + matches!(info.result, Ok(StepOutcome::Contention)), + "a held Serial step must report Contention to the surplus worker" + ); + assert_eq!(info.name, "", "the contended-step name is reported"); + assert_eq!( + runs.load(Ordering::Relaxed), + 0, + "the Serial step body must NOT run while another worker holds its mutex" + ); + + drop(guard); + } + /// A sticky-owned source driven through `run_worker_loop` completes and the /// loop exits even though the cached `sticky_live` flag (not a per-iteration /// `live.contains` scan) gates the fast path. Exercises the S1b-006 cache: @@ -966,6 +1173,8 @@ mod tests { &crate::runtime::scheduler::ChainOrderScheduler, None, 0, + None, + false, ); } @@ -1031,6 +1240,8 @@ mod tests { &crate::runtime::scheduler::ChainOrderScheduler, None, 0, + None, + false, ); // The source must have been called EXACTLY twice: call 1 = idle in the @@ -1176,6 +1387,8 @@ mod tests { &crate::runtime::scheduler::ChainOrderScheduler, None, 0, + None, + false, ); done.store(true, Ordering::SeqCst); @@ -1337,6 +1550,8 @@ mod tests { &crate::runtime::scheduler::ChainOrderScheduler, None, 0, + None, + false, ); done.store(true, Ordering::SeqCst); @@ -1418,6 +1633,7 @@ mod tests { true, None, 0, + None, ); let snap = stats.snapshot(); assert!( @@ -1447,6 +1663,7 @@ mod tests { false, None, 0, + None, ); assert!( stats_pool.snapshot().detached.is_empty(), @@ -1489,6 +1706,7 @@ mod tests { false, Some(&board), 0, + None, ); assert_eq!(board.read(0), (WorkerState::Running, Some(StepIdx(0)))); } @@ -1524,6 +1742,7 @@ mod tests { false, None, 0, + None, ) .unwrap(); assert!(matches!(info.result, Ok(StepOutcome::Finished))); diff --git a/crates/fgumi-pipeline-core/src/runtime/event_count.rs b/crates/fgumi-pipeline-core/src/runtime/event_count.rs new file mode 100644 index 000000000..af4f2b140 --- /dev/null +++ b/crates/fgumi-pipeline-core/src/runtime/event_count.rs @@ -0,0 +1,491 @@ +//! `PoolEventCount` — a Vyukov-style event-count for parking idle pool workers. +//! +//! # Why this exists +//! +//! The chain-builder pool is *replicated round-robin polling*: when a worker's +//! whole dispatch pass does no work it exponential-backoff *sleeps*, and there +//! is no producer→consumer wakeup. On an oversubscribed pipeline (more workers +//! than parallel work) the surplus workers churn poll→sleep→poll forever, ~93% +//! idle, and that churn measurably slows the productive workers. Parking those +//! workers on a condition variable removes the churn — but a bare condvar has a +//! lost-wakeup race: a producer that publishes work *after* an idle worker +//! checked its queues (saw empty) but *before* it blocks would wake nobody, and +//! the worker sleeps on work that is already present. +//! +//! An **event-count** closes that race. It is a generation counter plus a +//! waiter count, used in the canonical two-phase protocol: +//! +//! - **Worker (about to idle):** [`prepare_wait`](PoolEventCount::prepare_wait) +//! (register as a waiter, `SeqCst` fence, snapshot the generation) → re-check +//! the *real* condition (its own queues) → if still no work, +//! [`wait`](PoolEventCount::wait) (block until the generation moves or the +//! deadline elapses). +//! - **Producer (just published work):** [`notify_one`](PoolEventCount::notify_one) +//! (`SeqCst` fence, then read the waiter count; if any, bump the generation +//! under the lock and wake one). +//! +//! # Correctness — one rule +//! +//! The design reduces to the C++20 `SeqCst`-fence rule: if a `SeqCst` fence X is +//! sequenced after a write A on one thread, and a `SeqCst` fence Y is sequenced +//! before a read B on another, and X precedes Y in the single total order of +//! `SeqCst` operations, then B observes A (or a later write). The two fences are +//! the one in `prepare_wait` (after `waiters += 1`, before the worker's +//! re-poll) and the one in `notify_one` (after the producer's publish, before +//! its `waiters` load). They are totally ordered, so exactly one holds: +//! +//! - **Worker's fence precedes producer's.** The producer's `waiters` load +//! observes the increment (≥ 1), takes the lock, bumps `gen`, notifies. The +//! worker is either not yet at its `gen` check (which runs under the same +//! lock, after the producer releases it, so it sees the bump and returns +//! without sleeping) or already in `wait_for` (which the notify wakes). +//! - **Producer's fence precedes worker's.** The worker's re-poll (sequenced +//! after its fence) observes the producer's published item and never waits. +//! +//! No interleaving parks on a poppable item. The item handoff's own +//! happens-before is unaffected (it stays with the queue's atomics); this type +//! only orders the *wakeup*. Notifying under the same lock the condvar waits on +//! is what gives the condvar a "token": a bump that lands between the worker's +//! `gen` snapshot and its `wait_for` is not lost, because `wait` re-reads `gen` +//! under the lock before blocking. +//! +//! # Cost +//! +//! Steady state, nobody parked: `notify_one` is one `SeqCst` fence plus one +//! read-shared load of `waiters` — no store, no shared-line RMW, no lock. The +//! `waiters` and `gen` atomics are cache-line isolated (128-byte stride, as the +//! sharded liveness counter already does) so an idling worker's +//! increment/decrement of `waiters` never false-shares with a producer's load. + +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering, fence}; +use std::time::Duration; + +use parking_lot::{Condvar, Mutex}; + +/// Cache-line stride used to isolate the hot atomics from each other and from +/// any neighbouring fields. 128 bytes covers the 64-byte line plus the +/// adjacent-line prefetch pairing on aarch64 and x86-64 (matches the padding +/// the sharded liveness counter uses). +const CACHE_LINE: usize = 128; + +#[repr(align(128))] +struct Padded(T); + +/// A Vyukov event-count over a `parking_lot` `Mutex`/`Condvar`. +/// +/// Shared by all pool workers of one pipeline via `Arc`. Only present when +/// `n_threads > 1`; the fused and scheduled single-thread paths pass `None` and +/// keep their existing sleep-backoff idle. +pub struct PoolEventCount { + /// Number of workers currently *armed or blocked* (between `prepare_wait` + /// and the matching `wait`/`cancel_wait`). Read by `notify_one` on the hot + /// path to skip the lock+wake when nobody is parked. Cache-line isolated. + waiters: Padded, + /// Generation counter, bumped only when notifying, only under `lock`. A + /// worker that snapshots `gen` in `prepare_wait` and finds it changed by + /// the time it would block knows work was published in its wait window and + /// does not block. Cache-line isolated. + generation: Padded, + /// Guards the `gen` bump and the condvar. A producer takes it only when + /// `waiters > 0`; a worker takes it only to block. + lock: Mutex<()>, + cvar: Condvar, + /// Total pool worker count, so `notify_one` can compute + /// `awake = n_pool − waiters` for the ceiling gate. Set once at + /// construction; never mutated. + n_pool: usize, + /// Concurrency ceiling: `notify_one` wakes a parked worker only while + /// `awake < ceiling`. Above it, a push wakes nobody — the already-awake + /// workers find the item on their next pass (they cannot park while their + /// passes find work), so parking *sticks* instead of the whole surplus + /// re-waking on every push (the per-`Progress` cascade). Default `n_pool` + /// (ungated: `awake ≤ n_pool` always holds, so every notify fires) — the + /// behaviour before a controller adjusts it. `notify_all` (cancel / edge + /// close / `Finished`) is NEVER gated: correctness must not depend on the + /// ceiling. Cache-line isolated from the hot `waiters`/`gen`. + ceiling: Padded, +} + +const _: () = assert!(std::mem::align_of::>() == CACHE_LINE); + +/// Opaque generation snapshot from [`PoolEventCount::prepare_wait`], handed back +/// to [`PoolEventCount::wait`]. Carrying it by value (rather than re-reading +/// inside `wait`) is what makes the "generation moved in my wait window" check +/// precise. +/// +/// Deliberately **not** `Copy`/`Clone`: a key is a single-use token that +/// balances exactly one `prepare_wait` increment. `wait` and `cancel_wait` each +/// consume it by value, so the type system forbids handing one key to both (or +/// to either twice) — which would decrement `waiters` more than once and can +/// wrap the counter, causing a later `notify_one` to skip an armed worker until +/// its timeout. +#[derive(Debug, PartialEq, Eq)] +pub struct WaitKey(u64); + +/// Outcome of a [`PoolEventCount::wait`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum WaitOutcome { + /// The generation had already moved by the time `wait` took the lock — work + /// was published in the arm→wait window; the caller must re-poll, not sleep. + Woken, + /// The condvar returned (a notify, a spurious wake, or the deadline). The + /// caller loops and re-polls regardless — spurious wakes are harmless. + Returned, +} + +impl PoolEventCount { + /// Construct for a pool of `n_pool` workers. The ceiling starts at `n_pool` + /// (ungated). Use [`Self::set_ceiling`] to gate wakeups. + #[must_use] + pub fn new(n_pool: usize) -> Self { + Self { + waiters: Padded(AtomicUsize::new(0)), + generation: Padded(AtomicU64::new(0)), + lock: Mutex::new(()), + cvar: Condvar::new(), + n_pool, + ceiling: Padded(AtomicUsize::new(n_pool)), + } + } + + /// Set the concurrency ceiling — the max number of workers `notify_one` + /// will keep awake. Clamped to `1..=n_pool` (a ceiling of 0 would wedge the + /// pool: no worker could ever be woken). A controller (edge-depth / + /// rejection driven) calls this to shed or admit surplus workers; with no + /// controller the default `n_pool` leaves every notify ungated. + #[inline] + pub fn set_ceiling(&self, ceiling: usize) { + self.ceiling.0.store(ceiling.clamp(1, self.n_pool.max(1)), Ordering::Relaxed); + } + + /// The current ceiling (for the controller / tests). + #[must_use] + #[inline] + pub fn ceiling(&self) -> usize { + self.ceiling.0.load(Ordering::Relaxed) + } + + /// Number of workers currently armed or blocked. Used by the ceiling + /// controller to compute `awake = n_pool − waiters`. + #[must_use] + #[inline] + pub fn waiters(&self) -> usize { + self.waiters.0.load(Ordering::Relaxed) + } + + /// Producer side: wake at most one parked worker if any is waiting. + /// + /// Hot path when nobody is parked: one `SeqCst` fence + one relaxed load, + /// then return. No store, no lock, no wake. The fence is the producer half + /// of the correctness rule and must be `SeqCst`; it is sequenced *after* the + /// caller's publish (the successful queue push) and *before* the `waiters` + /// load, so a worker that armed before the publish is seen here. + #[inline] + pub fn notify_one(&self) { + fence(Ordering::SeqCst); + let waiters = self.waiters.0.load(Ordering::Relaxed); + if waiters == 0 { + return; + } + // Ceiling gate: keep at most `ceiling` workers awake. `awake = n_pool − + // waiters`; if we are already at/above the ceiling, do not wake another + // — the awake workers will find this item on their next pass. This is + // what stops the per-`Progress` cascade from re-waking the whole surplus + // on an oversubscribed pipeline. Ungated when `ceiling == n_pool` + // (default): then `awake < n_pool` whenever `waiters >= 1`, so the + // branch never suppresses a wake. + let awake = self.n_pool.saturating_sub(waiters); + if awake >= self.ceiling.0.load(Ordering::Relaxed) { + return; + } + let _guard = self.lock.lock(); + self.generation.0.fetch_add(1, Ordering::Relaxed); + self.cvar.notify_one(); + } + + /// Producer side: wake *every* parked worker. Reserved for terminal / + /// broadcast events (edge close, `Finished`, cancel/error) where more than + /// one waiter may need to observe the transition. Same fence discipline as + /// [`Self::notify_one`]. + #[inline] + pub fn notify_all(&self) { + fence(Ordering::SeqCst); + if self.waiters.0.load(Ordering::Relaxed) == 0 { + return; + } + let _guard = self.lock.lock(); + self.generation.0.fetch_add(1, Ordering::Relaxed); + self.cvar.notify_all(); + } + + /// Worker side, phase 1: register as a waiter and snapshot the generation. + /// + /// The `SeqCst` fence is the worker half of the correctness rule: it is + /// sequenced *after* the `waiters` increment and *before* the caller's + /// re-poll of its real condition. The returned [`WaitKey`] must be passed to + /// exactly one of [`Self::wait`] or [`Self::cancel_wait`] to balance the + /// increment. + #[must_use] + #[inline] + pub fn prepare_wait(&self) -> WaitKey { + self.waiters.0.fetch_add(1, Ordering::Relaxed); + fence(Ordering::SeqCst); + WaitKey(self.generation.0.load(Ordering::Relaxed)) + } + + /// Worker side: abandon a prepared wait without blocking (the re-poll found + /// work, or the pipeline is shutting down). Balances the `prepare_wait` + /// increment. + /// + /// Takes the `WaitKey` **by value** to consume it: the token is single-use + /// (see [`WaitKey`]), so taking it by reference would let one key balance + /// the increment twice. That is why the `needless_pass_by_value` lint is + /// suppressed here rather than followed. + #[inline] + #[allow(clippy::needless_pass_by_value)] + pub fn cancel_wait(&self, _key: WaitKey) { + self.waiters.0.fetch_sub(1, Ordering::Relaxed); + } + + /// Worker side, phase 2: block until the generation moves past `key` or + /// `deadline` elapses. Balances the `prepare_wait` increment on return. + /// + /// Returns [`WaitOutcome::Woken`] if the generation had already advanced + /// when the lock was taken (a notify raced into the arm→wait window — do not + /// block, re-poll immediately), else [`WaitOutcome::Returned`] after the + /// condvar wait (notify, spurious, or timeout — re-poll anyway). + /// + /// Consumes the `WaitKey` (see [`cancel_wait`](Self::cancel_wait) for why it + /// is taken by value rather than by reference). + #[must_use] + #[allow(clippy::needless_pass_by_value)] + pub fn wait(&self, key: WaitKey, deadline: Duration) -> WaitOutcome { + let mut guard = self.lock.lock(); + // Under the lock, re-read `gen`. If it moved since `prepare_wait`, a + // producer bumped it (also under this lock) after we armed — the wakeup + // we would wait for has already happened. This is the condvar's "token". + if self.generation.0.load(Ordering::Relaxed) != key.0 { + drop(guard); + self.waiters.0.fetch_sub(1, Ordering::Relaxed); + return WaitOutcome::Woken; + } + // `wait_for` may wake spuriously; the caller re-polls regardless, so a + // single wait (not a loop) is correct here. + let _ = self.cvar.wait_for(&mut guard, deadline); + drop(guard); + self.waiters.0.fetch_sub(1, Ordering::Relaxed); + WaitOutcome::Returned + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::thread; + + use super::*; + + #[test] + fn new_has_no_waiters() { + let ec = PoolEventCount::new(32); + assert_eq!(ec.waiters(), 0); + } + + #[test] + fn prepare_then_cancel_balances_waiters() { + let ec = PoolEventCount::new(32); + let key = ec.prepare_wait(); + assert_eq!(ec.waiters(), 1); + ec.cancel_wait(key); + assert_eq!(ec.waiters(), 0); + } + + #[test] + fn notify_one_with_no_waiters_is_a_noop_and_takes_no_lock() { + // The hot-path guarantee: with nobody parked, notify_one must not take + // the lock. Prove it by holding the lock for the duration of the call — + // if notify_one tried to lock, it would deadlock (parking_lot mutex is + // not reentrant) and the test would hang; a clean return proves the + // lock was never attempted. + let ec = PoolEventCount::new(32); + let held = ec.lock.lock(); + ec.notify_one(); // must return without touching the lock + drop(held); + assert_eq!(ec.waiters(), 0); + } + + #[test] + fn wait_returns_woken_when_generation_already_moved() { + // Models the arm→(notify)→wait window: the worker armed, a producer + // bumped the generation before the worker reached `wait`. `wait` must + // see the moved generation and return Woken without blocking. + let ec = PoolEventCount::new(32); + let key = ec.prepare_wait(); + // Simulate a producer's notify landing in the window (bumps gen under + // the lock, as notify_one does). + { + let _g = ec.lock.lock(); + ec.generation.0.fetch_add(1, Ordering::Relaxed); + } + let outcome = ec.wait(key, Duration::from_secs(30)); + assert_eq!(outcome, WaitOutcome::Woken, "a moved generation must not block"); + assert_eq!(ec.waiters(), 0, "wait balances the prepare_wait increment"); + } + + #[test] + fn wait_times_out_when_no_notify() { + // With no notify, wait blocks until the deadline and returns Returned. + let ec = PoolEventCount::new(32); + let key = ec.prepare_wait(); + let start = std::time::Instant::now(); + let outcome = ec.wait(key, Duration::from_millis(30)); + assert!(start.elapsed() >= Duration::from_millis(25), "must block ~the deadline"); + assert_eq!(outcome, WaitOutcome::Returned); + assert_eq!(ec.waiters(), 0); + } + + #[test] + fn notify_one_wakes_a_blocked_waiter() { + let ec = Arc::new(PoolEventCount::new(32)); + let ec2 = Arc::clone(&ec); + let woke = Arc::new(AtomicUsize::new(0)); + let woke2 = Arc::clone(&woke); + + let h = thread::spawn(move || { + let key = ec2.prepare_wait(); + // No work found; block for a long deadline so only a notify frees it. + let _ = ec2.wait(key, Duration::from_secs(30)); + woke2.store(1, Ordering::SeqCst); + }); + + // Wait until the worker has armed, then notify it. + while ec.waiters() == 0 { + thread::yield_now(); + } + // Give the worker a moment to reach the blocking wait_for. + thread::sleep(Duration::from_millis(20)); + ec.notify_one(); + + h.join().expect("waiter thread joins"); + assert_eq!(woke.load(Ordering::SeqCst), 1, "the blocked waiter was woken by notify_one"); + assert_eq!(ec.waiters(), 0); + } + + #[test] + fn notify_all_wakes_every_waiter() { + let ec = Arc::new(PoolEventCount::new(32)); + let n = 8usize; + let woke = Arc::new(AtomicUsize::new(0)); + let mut handles = Vec::new(); + for _ in 0..n { + let ec2 = Arc::clone(&ec); + let woke2 = Arc::clone(&woke); + handles.push(thread::spawn(move || { + let key = ec2.prepare_wait(); + let _ = ec2.wait(key, Duration::from_secs(30)); + woke2.fetch_add(1, Ordering::SeqCst); + })); + } + while ec.waiters() < n { + thread::yield_now(); + } + thread::sleep(Duration::from_millis(20)); + ec.notify_all(); + for h in handles { + h.join().expect("waiter joins"); + } + assert_eq!(woke.load(Ordering::SeqCst), n, "notify_all woke every waiter"); + assert_eq!(ec.waiters(), 0); + } + + /// Stress: producers publishing concurrently with a worker arming/waiting + /// must never leave the worker blocked on a "published" item. Models the + /// lost-wakeup window many times. Not a proof (that needs loom), but catches + /// gross ordering regressions. + #[test] + fn stress_no_lost_wakeup() { + for _ in 0..2000 { + let ec = Arc::new(PoolEventCount::new(32)); + let ec_p = Arc::clone(&ec); + // A flag standing in for "work is available", published before notify. + let work = Arc::new(AtomicUsize::new(0)); + let work_p = Arc::clone(&work); + + let producer = thread::spawn(move || { + work_p.store(1, Ordering::SeqCst); // publish + ec_p.notify_one(); // then notify + }); + + // Worker: arm, re-check the real condition, block only if no work. + let deadline = Duration::from_millis(500); + let key = ec.prepare_wait(); + let outcome = if work.load(Ordering::SeqCst) == 1 { + ec.cancel_wait(key); + WaitOutcome::Woken + } else { + // Blocked: a healthy notify (or spurious wake) frees us + // promptly; a *lost* wakeup would instead burn the full + // deadline and time out. `WaitOutcome::Returned` conflates + // notify and timeout, so it cannot tell them apart on its own — + // time the wait and assert it completed well within the + // deadline, which a timeout (lost wakeup) never does. + let start = std::time::Instant::now(); + let outcome = ec.wait(key, deadline); + assert!( + start.elapsed() < deadline / 2, + "blocked wait must complete promptly on notify, not time out (lost wakeup)" + ); + outcome + }; + producer.join().expect("producer joins"); + // Either the re-check saw the work, or the blocked wait returned + // well within the deadline (asserted above). Confirm the published + // work is visible on the way out. + assert_eq!(work.load(Ordering::SeqCst), 1); + let _ = outcome; + assert_eq!(ec.waiters(), 0); + } + } + + #[test] + fn set_ceiling_clamps_to_pool_and_floor_one() { + let ec = PoolEventCount::new(8); + assert_eq!(ec.ceiling(), 8, "default ceiling is the pool size (ungated)"); + ec.set_ceiling(3); + assert_eq!(ec.ceiling(), 3); + ec.set_ceiling(0); + assert_eq!(ec.ceiling(), 1, "0 clamps to 1 — a 0 ceiling would wedge the pool"); + ec.set_ceiling(100); + assert_eq!(ec.ceiling(), 8, "above n_pool clamps to n_pool"); + } + + #[test] + fn notify_one_suppressed_at_ceiling_but_notify_all_still_wakes() { + // n_pool = 4, ceiling = 1. One waiter parks → awake = 4 − 1 = 3 ≥ 1, so + // `notify_one` must NOT wake it (we are already at/above the ceiling). + // `notify_all` (ungated) must still wake it. + let ec = Arc::new(PoolEventCount::new(4)); + ec.set_ceiling(1); + let ec2 = Arc::clone(&ec); + let woke = Arc::new(AtomicUsize::new(0)); + let woke2 = Arc::clone(&woke); + let h = thread::spawn(move || { + let key = ec2.prepare_wait(); + let _ = ec2.wait(key, Duration::from_secs(30)); + woke2.store(1, Ordering::SeqCst); + }); + while ec.waiters() == 0 { + thread::yield_now(); + } + thread::sleep(Duration::from_millis(20)); + // Gated: awake (3) >= ceiling (1) → no wake. + ec.notify_one(); + thread::sleep(Duration::from_millis(30)); + assert_eq!(woke.load(Ordering::SeqCst), 0, "notify_one must be suppressed at the ceiling"); + // Ungated broadcast frees it. + ec.notify_all(); + h.join().expect("waiter joins"); + assert_eq!(woke.load(Ordering::SeqCst), 1, "notify_all is never gated"); + } +} diff --git a/crates/fgumi-pipeline-core/src/runtime/mod.rs b/crates/fgumi-pipeline-core/src/runtime/mod.rs index bce0ca709..c4ad70f41 100644 --- a/crates/fgumi-pipeline-core/src/runtime/mod.rs +++ b/crates/fgumi-pipeline-core/src/runtime/mod.rs @@ -5,6 +5,7 @@ pub mod contexts; pub mod detached; pub mod drain; pub mod driver; +pub mod event_count; pub mod fused; pub mod live; pub mod metrics; @@ -23,6 +24,7 @@ pub use detached::{ }; pub use drain::StepDrainCounter; pub use driver::run_worker_loop; +pub use event_count::{PoolEventCount, WaitKey, WaitOutcome}; pub use fused::{is_fusible_chain, run_fused_single_thread, should_fuse_single_thread}; pub use live::LiveSteps; pub use pool::{assign_exclusive_owners, assign_sticky_owners}; diff --git a/crates/fgumi-pipeline-core/src/runtime/worker_core.rs b/crates/fgumi-pipeline-core/src/runtime/worker_core.rs index 2274419a0..204b10f48 100644 --- a/crates/fgumi-pipeline-core/src/runtime/worker_core.rs +++ b/crates/fgumi-pipeline-core/src/runtime/worker_core.rs @@ -6,7 +6,17 @@ use crate::topology::StepIdx; // Pool worker idle bounds: a pool worker is never unparked, so it sleeps and can // ramp to a coarse cap without hurting wake latency. -const SLEEP_INITIAL_US: u64 = 1; +// +// The initial (also the post-progress reset floor) is 20µs, NOT 1µs. A surplus +// worker in an oversubscribed pool makes the occasional lucky pop, which resets +// its backoff; a 1µs floor then puts it back on the steep part of the ramp, +// waking ~every microsecond to run a full empty poll pass (touching every shared +// queue and probing every `Serial` step) — the churn Layer 0 targets. Wake +// latency for a worker that just made progress is unaffected: it did work, so it +// does not sleep at all this pass. The floor only blunts how fast an idle +// surplus worker re-ramps after a stray success. 20µs is well below any step's +// per-item service time yet coarse enough to stop the per-microsecond re-poll. +const SLEEP_INITIAL_US: u64 = 20; const SLEEP_MAX_US: u64 = 50_000; // 50 milliseconds // Dedicated-driver idle bounds: a driver drives a small step subset off the pool @@ -154,6 +164,15 @@ impl WorkerCore { self.backoff_us = self.backoff_us.saturating_mul(2).min(self.policy().max_us()); } + /// The current backoff as a `Duration`, used as the event-count `wait` + /// timeout (the self-heal bound on a missed wakeup). Same ramp as + /// `sleep_backoff` would sleep, so a parked worker's timeout grows on + /// repeated idle just as the spin-sleep did. + #[must_use] + pub fn backoff_deadline(&self) -> Duration { + Duration::from_micros(self.backoff_us) + } + #[cfg(test)] fn current_backoff_us(&self) -> u64 { self.backoff_us diff --git a/crates/fgumi-pipeline-core/src/signal.rs b/crates/fgumi-pipeline-core/src/signal.rs index eda14a32e..8a5c99f57 100644 --- a/crates/fgumi-pipeline-core/src/signal.rs +++ b/crates/fgumi-pipeline-core/src/signal.rs @@ -98,6 +98,17 @@ impl std::error::Error for PipelineError { pub struct PipelineSignal { state: AtomicU8, payload: OnceLock, + /// Late-bound handle to the pool's event-count parker, set by + /// `Pipeline::run` before any worker spawns (only when `n_threads > 1`; the + /// single-thread paths leave it empty). On a terminal transition + /// (`cancel`/`record_error`) the signal wakes every parked worker so they + /// observe `is_done()` promptly instead of sleeping to their wait timeout. + /// + /// Held as a `Weak` to avoid an `Arc` cycle: the event-count does not need + /// to keep the signal alive, and the signal must not keep the event-count + /// alive past the run. Workers hold the strong `Arc` for the run's duration, + /// so the upgrade always succeeds while a worker could be parked. + event_count: OnceLock>, } impl PipelineSignal { @@ -130,6 +141,25 @@ impl PipelineSignal { self.state.load(AtomicOrdering::Relaxed) == STATE_CANCELLED } + /// Bind the pool's event-count parker so terminal transitions wake parked + /// workers. Called once by `Pipeline::run` before spawning workers, only + /// when `n_threads > 1`. Idempotent-ish: a second call is silently dropped + /// (`OnceLock::set`), which never happens in practice (one `run` per + /// signal). Stored as a `Weak` to avoid an `Arc` cycle. + pub(crate) fn bind_event_count(&self, ec: &Arc) { + let _ = self.event_count.set(Arc::downgrade(ec)); + } + + /// Wake every parked worker after a terminal transition, so they observe + /// `is_done()` rather than sleeping to their wait timeout. No-op when no + /// event-count is bound (single-thread paths) or it has already been + /// dropped (all workers joined — nobody to wake). + fn wake_parked_workers(&self) { + if let Some(ec) = self.event_count.get().and_then(std::sync::Weak::upgrade) { + ec.notify_all(); + } + } + /// First writer wins; later writers are silently dropped (the /// `compare_exchange` rejects the state transition and the /// `OnceLock::set` rejects the payload write). @@ -151,6 +181,13 @@ impl PipelineSignal { { let _ = self.payload.set(err); } + // Wake any parked workers unconditionally: whether or not this call won + // the state CAS, the pipeline is now terminal and parked workers must + // observe `is_done()` rather than sleep to their wait timeout. A + // redundant `notify_all` (if a prior writer already woke them) is + // harmless. Sequenced AFTER the state store so a woken worker's + // `is_done()` re-check sees the terminal state. + self.wake_parked_workers(); } /// First writer wins; later writers (including a `record_error` after @@ -178,6 +215,9 @@ impl PipelineSignal { { let _ = self.payload.set(PipelineError::Cancelled); } + // See `record_error`: wake parked workers unconditionally after the + // terminal state store so they observe `is_done()` promptly. + self.wake_parked_workers(); } /// Read the recorded error payload, if any. From 223ee29974435f22bfbf6d04c519e814f37d242f Mon Sep 17 00:00:00 2001 From: Nils Homer Date: Sat, 12 Sep 2026 18:25:21 -0700 Subject: [PATCH 2/2] feat(correct): pin the serial GroupByQueryname to one worker MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `GroupByQueryname` is the throughput ceiling of the `correct` chain — the sole `Serial` feeder of the parallel `correct` step. Left as `Affinity::None` it is probed by every one of N workers on every idle pass; the engine's own bottleneck verdict flags a contended Serial step as an affinity/fuse/detach candidate. Give the step an optional per-instance affinity (default `Affinity::None`, so every other chain that uses it — align, dedup, filter-by-template — is unchanged) and, in `add_correct`, pin it to `Affinity::Worker(1)` (worker 1, not 0, which hosts the sticky reader source; clamped so `--threads 1` maps to worker 0 rather than requesting a non-existent worker). Surplus workers then `Skip` it entirely instead of thrashing its mutex. Pairs with the pool-worker parking in the engine: pinning stops the contention, parking stops the idle churn. Output is unchanged (pinning only selects which worker runs the serial step). --- src/lib/pipeline/chains/builder.rs | 15 ++++++++++++- src/lib/pipeline/steps/group/queryname.rs | 26 +++++++++++++++++++++++ 2 files changed, 40 insertions(+), 1 deletion(-) diff --git a/src/lib/pipeline/chains/builder.rs b/src/lib/pipeline/chains/builder.rs index df5e39cf4..fef198e2f 100644 --- a/src/lib/pipeline/chains/builder.rs +++ b/src/lib/pipeline/chains/builder.rs @@ -2286,7 +2286,20 @@ impl<'a> ChainBuilder<'a> { let tail = if self.chain_tail_kind == ChainTailKind::BamTemplateBatch { tail } else { - self.pipeline.append_step(GroupByQueryname::new(self.tuning.per_step_byte_limit), tail) + // Layer 1 (oversubscription): `GroupByQueryname` is the throughput + // ceiling of the `correct` chain — it is the sole `Serial` feeder of + // the parallel `correct` step. Left as `Affinity::None` it is probed + // by every one of N workers on every idle pass; pin it to a single + // worker so the surplus `Skip` it entirely. Worker 1 (not 0, which + // hosts the sticky reader source), clamped so `--threads 1` maps to + // worker 0 rather than requesting a non-existent worker (which would + // panic at run start). + use crate::pipeline::core::step::Affinity; + let pin = Affinity::Worker(1.min(num_threads.saturating_sub(1))); + self.pipeline.append_step( + GroupByQueryname::new(self.tuning.per_step_byte_limit).with_affinity(pin), + tail, + ) }; // Dispatch on rejects presence: either a 2-output or a 1-output step. diff --git a/src/lib/pipeline/steps/group/queryname.rs b/src/lib/pipeline/steps/group/queryname.rs index 160d050d5..b132b8743 100644 --- a/src/lib/pipeline/steps/group/queryname.rs +++ b/src/lib/pipeline/steps/group/queryname.rs @@ -62,6 +62,14 @@ pub struct GroupByQueryname { target_batch_count: usize, output_byte_limit: u64, name: &'static str, + /// Scheduling affinity for this `Serial` step. Default `Affinity::None` + /// (every worker may attempt it) preserves the behaviour of every chain + /// that does not opt in. A chain where this step is the throughput ceiling + /// (e.g. `correct`, whose parallel `correct` step is fed only by this + /// grouper) can pin it to a single worker via [`Self::with_affinity`] so + /// surplus workers `Skip` it entirely instead of probing its mutex every + /// idle pass (Layer 1 of the oversubscription fix). + affinity: crate::pipeline::core::step::Affinity, } impl GroupByQueryname { @@ -84,9 +92,23 @@ impl GroupByQueryname { target_batch_count: target_batch_count.max(1), output_byte_limit, name: "GroupByQueryname", + affinity: crate::pipeline::core::step::Affinity::None, } } + /// Pin this `Serial` grouper to a single worker (Layer 1 oversubscription + /// fix). With a non-`None` affinity the step is a build-time `Skip` on every + /// other worker, so no surplus worker ever probes or takes its mutex — use + /// when this grouper is a chain's throughput ceiling. The caller must pass a + /// worker index `< n_threads` (`Affinity::Worker(k)` with `k >= n_threads` + /// panics at run start); prefer a non-zero `k` so it does not share worker 0 + /// with a sticky reader source. + #[must_use] + pub fn with_affinity(mut self, affinity: crate::pipeline::core::step::Affinity) -> Self { + self.affinity = affinity; + self + } + /// Flush the run-in-progress into the accumulator (if non-empty). The /// `current_name` buffer is retained (only the run-active flag is cleared) /// so its allocation is reused by the next template boundary. @@ -170,6 +192,10 @@ impl Step for GroupByQueryname { } } + fn affinity(&self) -> crate::pipeline::core::step::Affinity { + self.affinity + } + fn try_run(&mut self, ctx: &mut StepCtx<'_, Self>) -> io::Result { // 1. Drain held slot first. if let Some(unpushed) = self.held.take() {