diff --git a/crates/fgumi-pipeline-core/src/builder.rs b/crates/fgumi-pipeline-core/src/builder.rs index d7b8cb455..c94f98fdd 100644 --- a/crates/fgumi-pipeline-core/src/builder.rs +++ b/crates/fgumi-pipeline-core/src/builder.rs @@ -1259,6 +1259,7 @@ impl Pipeline { config.instrumentation, config.telemetry.is_some(), )); + scheduler.bind(&contexts.bounded_queues); // 2-pre-monitor invariant: if the deadlock monitor will be armed, every // output transport must be ByteBounded so `in_flight_bytes` can see a diff --git a/crates/fgumi-pipeline-core/src/header.rs b/crates/fgumi-pipeline-core/src/header.rs index 7109d33b7..a692befdd 100644 --- a/crates/fgumi-pipeline-core/src/header.rs +++ b/crates/fgumi-pipeline-core/src/header.rs @@ -3,7 +3,7 @@ //! `HeaderHandle` lets a downstream sink open its output file before the //! full output header is known, and consume the header at first record //! time once an upstream step has produced it. The primary motivating -//! consumer is the writer downstream of `AlignAndMergeStep`: the +//! consumer is the writer downstream of the align stage: the //! aligner's `@PG` (and any `@RG`/`@CO` lines it adds) are runtime //! contributions that aren't available at `Pipeline::build` time. //! diff --git a/crates/fgumi-pipeline-core/src/runtime/driver.rs b/crates/fgumi-pipeline-core/src/runtime/driver.rs index 80fe91c8c..e0584d716 100644 --- a/crates/fgumi-pipeline-core/src/runtime/driver.rs +++ b/crates/fgumi-pipeline-core/src/runtime/driver.rs @@ -372,6 +372,27 @@ struct RoundRobinOutcome { removed_sticky_owner: bool, } +/// For [`WalkDirection::RefillThenReverse`], how many of the leading live +/// steps (`order`, in chain order) are at or before the refill step and so walk +/// forward; `0` for the other directions. +fn refill_split(walk: WalkDirection, order: &[StepIdx]) -> usize { + match walk { + WalkDirection::RefillThenReverse(through) => order.partition_point(|s| s.0 <= through.0), + WalkDirection::Forward | WalkDirection::Reverse => 0, + } +} + +/// The live-step position the `i`-th visit of a pass lands on, for `n` live +/// steps of which the first `split` walk forward (see [`refill_split`]). +fn walk_position(walk: WalkDirection, i: usize, n: usize, split: usize) -> usize { + match walk { + WalkDirection::Forward => i, + WalkDirection::Reverse => n - 1 - i, + WalkDirection::RefillThenReverse(_) if i < split => i, + WalkDirection::RefillThenReverse(_) => n - 1 - (i - split), + } +} + /// One pass of the round-robin dispatch over this worker's live steps, in /// chain order. Finished steps are removed from `live` at end-of-pass (deferred /// so the in-progress walk over `live.order()` is not mutated underneath it). @@ -405,6 +426,7 @@ fn round_robin_dispatch( // safe and avoids reorder-under-iteration. let mut finished: Vec = Vec::new(); let n = live.len(); + let split = refill_split(walk, live.order()); for i in 0..n { if signal.is_done() { break; @@ -414,10 +436,7 @@ fn round_robin_dispatch( // `Reverse` = downstream-first (favour draining buffered work before // producing more). Everything else — skip-on-contention, the sticky // source/sink fast-path, Finished handling — is direction-agnostic. - let pos = match walk { - WalkDirection::Forward => i, - WalkDirection::Reverse => n - 1 - i, - }; + let pos = walk_position(walk, i, n, split); let step_idx = live.order()[pos]; let entry = &mut entries[step_idx.0]; let mut mark_skip = false; @@ -444,13 +463,22 @@ fn round_robin_dispatch( // Restart the walk at the top so the highest-priority step runs // again — EXCEPT for this worker's sticky owner. That step just // had its dedicated burst in phase 1 of the worker loop, so a - // priority restart here only re-runs it. Worse, `Progress` also - // means "held an item", so a sticky step sitting on a full output - // reports it forever: breaking the pass on that outcome means the - // downstream step that would drain the output is never reached - // and the pair livelocks at one worker. Walking past the sticky - // owner is what turns the bounded burst into real forward - // progress. + // priority restart here only re-runs it. Worse, a sticky step + // that reports `Progress` while sitting on a full output (against + // the `StepOutcome::Progress` contract) would report it forever: + // breaking the pass on that outcome means the downstream step + // that would drain the output is never reached and the pair + // livelocks at one worker. Walking past the sticky owner is what + // turns the bounded burst into real forward progress. + // + // A non-sticky step gets the restart, and that is safe under + // every walk (including the forward refill prefix) because of + // the held-slot convention its steps follow: a *new* hold reports + // `Progress` once, while a retry that is still rejected reports + // `NoProgress`/`Contention`. The restart therefore costs one + // extra visit, and that visit's idle outcome lets the walk reach + // the consumer that drains the output + // (`refill_walk_reaches_the_consumer_of_a_holding_step`). restart_priority = sticky_owner != Some(step_idx); } // Nothing to do this call. The step terminates by returning @@ -696,6 +724,30 @@ mod tests { use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use super::*; + + /// The visit order for each walk over five live steps `[0, 2, 3, 5, 7]` + /// (chain order, with gaps where steps finished). + #[rstest::rstest] + #[case::forward(WalkDirection::Forward, &[0, 2, 3, 5, 7])] + #[case::reverse(WalkDirection::Reverse, &[7, 5, 3, 2, 0])] + #[case::refill_through_live_step( + WalkDirection::RefillThenReverse(StepIdx(3)), + &[0, 2, 3, 7, 5] + )] + #[case::refill_through_finished_step( + WalkDirection::RefillThenReverse(StepIdx(4)), + &[0, 2, 3, 7, 5] + )] + #[case::refill_through_first(WalkDirection::RefillThenReverse(StepIdx(0)), &[0, 7, 5, 3, 2])] + #[case::refill_through_last(WalkDirection::RefillThenReverse(StepIdx(7)), &[0, 2, 3, 5, 7])] + #[case::refill_before_all(WalkDirection::RefillThenReverse(StepIdx(9)), &[0, 2, 3, 5, 7])] + fn walk_visits_live_steps_in_order(#[case] walk: WalkDirection, #[case] expected: &[usize]) { + let order: Vec = [0, 2, 3, 5, 7].into_iter().map(StepIdx).collect(); + let split = refill_split(walk, &order); + let visited: Vec = + (0..order.len()).map(|i| order[walk_position(walk, i, order.len(), split)].0).collect(); + assert_eq!(visited, expected); + } use crate::erased::{ErasedStep, TypedStep}; use crate::handles::BranchInputHandle; use crate::outputs::Single; @@ -1398,8 +1450,9 @@ mod tests { // ── Sticky forward-progress ───────────────────────────────────────────── // - // `StepOutcome::Progress` means "pushed OR HELD an item", so a sticky step - // whose output is full keeps reporting `Progress` while moving nothing. Two + // The source below breaks the `StepOutcome::Progress` contract on purpose: + // it keeps reporting `Progress` while its output is full and it moves + // nothing. Even so, a sticky pair must not wedge. Two // things must hold for the pair below to make progress at ONE worker: the // sticky burst is bounded, and round-robin does not restart the walk on the // sticky owner's `Progress` (which would break the pass before the sink is @@ -1407,9 +1460,10 @@ mod tests { // its_draining_consumer`. /// Sticky source over a capacity-1 output. Emits `remaining` items; when the - /// transport rejects a push it *holds* the item and reports `Progress` — the - /// contract's "pushed or held" case, and the outcome that makes an unbounded - /// sticky loop spin forever. + /// transport rejects a push it *holds* the item and reports `Progress`, even + /// on a retry that is still rejected (which the contract says should be + /// `NoProgress`). That is the outcome that makes an unbounded sticky loop + /// spin forever. struct StickyHoldingSource { remaining: u32, held: Option, @@ -1573,6 +1627,139 @@ mod tests { ); } + /// Non-sticky source over a capacity-1 output that follows the held-slot + /// convention every production step uses: a *new* hold reports `Progress` + /// (it claimed an item), while a retry that is still rejected reports + /// `NoProgress` (it moved nothing). + struct HoldingSource { + remaining: u32, + held: Option, + } + impl Step for HoldingSource { + type Input = (); + type Outputs = Single; + fn profile(&self) -> StepProfile { + StepProfile { + name: "HoldingSource", + kind: StepKind::Exclusive, + sticky: false, + output_queues: vec![QueueSpec::CountBounded { capacity: 1 }], + branch_ordering: vec![BranchOrdering::None], + } + } + fn try_run(&mut self, ctx: &mut StepCtx<'_, Self>) -> io::Result { + if let Some(item) = self.held.take() { + return match ctx.outputs.push(item) { + Ok(()) => Ok(StepOutcome::Progress), + Err(unpushed) => { + self.held = Some(unpushed.into_item()); + Ok(StepOutcome::NoProgress) + } + }; + } + if self.remaining == 0 { + return Ok(StepOutcome::Finished); + } + let item = self.remaining; + self.remaining -= 1; + if let Err(unpushed) = ctx.outputs.push(item) { + self.held = Some(unpushed.into_item()); + } + Ok(StepOutcome::Progress) + } + } + + /// Under a permanently raised refill signal, one worker walking a non-sticky + /// source that holds on a full output must still reach the sink that drains + /// it, whether the sink falls after the refill prefix (reverse part) or + /// inside it (forward part). The priority restart after the source's new + /// hold costs one extra visit, whose still-held retry reports `NoProgress`, + /// so the walk continues to the sink instead of restarting again. + #[rstest::rstest] + #[case::sink_after_refill_prefix(StepIdx(0))] + #[case::sink_inside_refill_prefix(StepIdx(1))] + fn refill_walk_reaches_the_consumer_of_a_holding_step(#[case] refill_through: StepIdx) { + const N_ITEMS: u32 = 4; + + let mut graph = ChainGraph::new(); + let src = graph.register_step("HoldingSource", 1); + let sink = graph.register_step("DrainingSink", 0); + graph.wire(src, BranchIdx(0), sink); + + let received = Arc::new(Mutex::new(Vec::new())); + let steps: Vec> = vec![ + Box::new(TypedStep::new(HoldingSource { remaining: N_ITEMS, held: None })), + Box::new(TypedStep::new(DrainingSink { received: Arc::clone(&received) })), + ]; + let contexts = Arc::new(build_chain_contexts( + &steps, + &graph, + crate::builder::InstrumentationLevel::Off, + false, + )); + + let mut entries: Vec = + steps.into_iter().map(|step| WorkerStepEntry::Exclusive { step }).collect(); + let drain_counters = vec![StepDrainCounter::new(1), StepDrainCounter::new(1)]; + let signal = PipelineSignal::new(); + + // Never bound, so the read-ahead cap is absent and the refill walk + // applies for the whole run. + let scheduler = crate::runtime::scheduler::RefillDrainScheduler::new( + Arc::new(AtomicBool::new(true)), + refill_through, + BranchIdx(0), + u64::MAX, + ); + assert_eq!( + crate::runtime::scheduler::Scheduler::walk(&scheduler), + WalkDirection::RefillThenReverse(refill_through) + ); + + // Watchdog: a walk that never reaches the sink spins forever. + let done = Arc::new(AtomicBool::new(false)); + { + let done = Arc::clone(&done); + std::thread::spawn(move || { + for _ in 0..400 { + std::thread::sleep(std::time::Duration::from_millis(25)); + if done.load(Ordering::SeqCst) { + return; + } + } + eprintln!( + "refill_walk_reaches_the_consumer_of_a_holding_step: WEDGED — the refill \ + walk never reaches the holding step's consumer" + ); + std::process::abort(); + }); + } + + let mut worker = WorkerCore::new(0, None, None); + run_worker_loop( + &mut worker, + &mut entries, + &contexts, + &drain_counters, + &signal, + None, + &crate::liveness::LivenessCounter::new(1), + &scheduler, + None, + 0, + None, + false, + ); + done.store(true, Ordering::SeqCst); + + assert_eq!( + *received.lock(), + (1..=N_ITEMS).rev().collect::>(), + "every item must reach the sink, in emission order" + ); + assert!(!signal.is_done(), "clean completion, no error"); + } + /// A step that spends a measurable, nonzero span inside `try_run` so its /// dispatch busy-time rounds above 0 ns. #[derive(Clone)] diff --git a/crates/fgumi-pipeline-core/src/runtime/mod.rs b/crates/fgumi-pipeline-core/src/runtime/mod.rs index c4ad70f41..c7a25bda3 100644 --- a/crates/fgumi-pipeline-core/src/runtime/mod.rs +++ b/crates/fgumi-pipeline-core/src/runtime/mod.rs @@ -28,7 +28,9 @@ 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}; -pub use scheduler::{ChainOrderScheduler, DrainFirstScheduler, Scheduler, WalkDirection}; +pub use scheduler::{ + ChainOrderScheduler, DrainFirstScheduler, RefillDrainScheduler, Scheduler, WalkDirection, +}; pub use stats::{PipelineStats, StatsSnapshot, StepStatsSnapshot}; pub use storage::{WorkerStepEntry, build_worker_storage}; pub use worker_core::{BackoffPolicy, WorkerCore, WorkerRole}; diff --git a/crates/fgumi-pipeline-core/src/runtime/scheduler.rs b/crates/fgumi-pipeline-core/src/runtime/scheduler.rs index 0fde6eb00..874d1244a 100644 --- a/crates/fgumi-pipeline-core/src/runtime/scheduler.rs +++ b/crates/fgumi-pipeline-core/src/runtime/scheduler.rs @@ -37,6 +37,10 @@ pub enum WalkDirection { Forward, /// Reverse chain order (downstream → upstream): favour draining. Reverse, + /// Chain order through the given step, then reverse chain order for the + /// steps after it: refill a starved stage's input first, while everything + /// downstream of it keeps draining before it produces more. + RefillThenReverse(crate::topology::StepIdx), } /// A per-worker dispatch-order policy. Selected per pipeline via @@ -49,6 +53,11 @@ pub trait Scheduler: Send + Sync + std::fmt::Debug { /// Human-readable name for diagnostics / `--pipeline-stats`. fn name(&self) -> &'static str; + + /// Called once per run with the run's byte-bounded queues, after they exist + /// and before any worker dispatches. A policy that watches a queue resolves + /// its handle here. The shipped stateless policies ignore it. + fn bind(&self, _queues: &[crate::runtime::contexts::RegisteredQueue]) {} } /// Default upstream-first (chain-order) walk. Preserves the historical dispatch @@ -83,6 +92,125 @@ impl Scheduler for DrainFirstScheduler { } } +/// Downstream-first dispatch, except while a stage has raised a shared +/// *refill* signal: then the steps up to and including that stage walk +/// upstream-first, and the rest stay downstream-first. +/// +/// [`DrainFirstScheduler`] visits the source and early steps last, so while a +/// heavy downstream step always has work, nothing upstream runs and the next +/// batch of input is never read ahead; the stage later waits on just-in-time +/// decoding. A stage that needs input raises the signal, and while it is raised +/// its feeding steps (`refill_through` and everything before it) come first in +/// the walk. Only those steps move: turning the whole walk upstream-first would +/// also let the stage's own heavy consumers run ahead of the steps draining +/// older work, which grows the working set and costs CPU. The signal is read +/// once per worker pass with a relaxed load. +/// +/// Read-ahead is capped: the refill walk applies only while the watched output +/// branch of `refill_through` (the queue feeding the stage) holds less than +/// `cap_bytes`, so the stage gets about one batch of work ready rather than +/// every upstream queue filling to its limit. Naming the branch keeps the cap +/// on the right queue when `refill_through` fans out (e.g. a rejects branch). +/// +/// One instance serves one [`Pipeline::run`](crate::builder::Pipeline::run): it +/// names that pipeline's step and branch, and binds that run's queue once. +/// Build a fresh one per run rather than reusing a cloned `PipelineConfig`. +pub struct RefillDrainScheduler { + refill: std::sync::Arc, + refill_through: crate::topology::StepIdx, + refill_branch: crate::topology::BranchIdx, + cap_bytes: u64, + /// The watched output queue, resolved in [`Scheduler::bind`]. + queue: std::sync::OnceLock>, +} + +impl std::fmt::Debug for RefillDrainScheduler { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("RefillDrainScheduler") + .field("refill_through", &self.refill_through) + .field("refill_branch", &self.refill_branch) + .field("cap_bytes", &self.cap_bytes) + .field("bound", &self.queue.get().is_some()) + .finish_non_exhaustive() + } +} + +impl RefillDrainScheduler { + /// A scheduler that walks `refill_through` and the steps before it + /// upstream-first while `refill` is set and `refill_through`'s output + /// branch `refill_branch` holds less than `cap_bytes`. + #[must_use] + pub fn new( + refill: std::sync::Arc, + refill_through: crate::topology::StepIdx, + refill_branch: crate::topology::BranchIdx, + cap_bytes: u64, + ) -> Self { + Self { refill, refill_through, refill_branch, cap_bytes, queue: std::sync::OnceLock::new() } + } + + /// Whether the refill walk applies now. + fn refilling(&self) -> bool { + self.refill.load(std::sync::atomic::Ordering::Relaxed) + && self.queue.get().is_none_or(|q| q.current_bytes() < self.cap_bytes) + } +} + +impl Scheduler for RefillDrainScheduler { + #[inline] + fn walk(&self) -> WalkDirection { + if self.refilling() { + WalkDirection::RefillThenReverse(self.refill_through) + } else { + WalkDirection::Reverse + } + } + + fn name(&self) -> &'static str { + "refill-drain" + } + + /// Resolve the watched output branch's queue. If it is not byte-bounded + /// the cap cannot be read and the refill walk follows the signal alone. + /// + /// Binding a second run's (different) queue is a contract violation — the + /// `OnceLock` would keep watching the first run's queue — so it trips a + /// `debug_assert!` and logs a warning in release. + fn bind(&self, queues: &[crate::runtime::contexts::RegisteredQueue]) { + if let Some(queue) = queues + .iter() + .find(|q| q.producer_step == self.refill_through && q.branch == self.refill_branch) + { + if let Err(handle) = self.queue.set(std::sync::Arc::clone(&queue.handle)) { + let same_queue = self.queue.get().is_some_and(|bound| { + std::ptr::addr_eq( + std::sync::Arc::as_ptr(bound), + std::sync::Arc::as_ptr(&handle), + ) + }); + debug_assert!( + same_queue, + "RefillDrainScheduler bound to a second run's queue; one instance serves \ + one Pipeline::run" + ); + if !same_queue { + log::warn!( + "refill-drain: scheduler reused across pipeline runs; its read-ahead \ + cap keeps watching the first run's queue" + ); + } + } + } else { + log::debug!( + "refill-drain: no byte-bounded queue for step {:?} branch {:?}; \ + read-ahead is uncapped and follows the refill signal alone", + self.refill_through, + self.refill_branch, + ); + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -99,8 +227,152 @@ mod tests { assert_eq!(DrainFirstScheduler.name(), "drain-first"); } - // The forward/reverse position→step arithmetic these policies drive is - // exercised end-to-end against the real dispatcher in driver.rs - // (`driver_round_robins_all_live_before_parking`); a standalone test here - // would only re-derive the formula, not the driver's actual mapping. + // The position→step arithmetic these policies drive is tested where it + // lives, in driver.rs: `walk_visits_live_steps_in_order` pins every walk's + // visit order, and `driver_round_robins_all_live_before_parking` / + // `refill_walk_reaches_the_consumer_of_a_holding_step` run it end-to-end + // against the real dispatcher. + + #[test] + fn refill_drain_walks_upstream_only_while_the_signal_is_raised() { + use std::sync::Arc; + use std::sync::atomic::{AtomicBool, Ordering}; + let refill = Arc::new(AtomicBool::new(false)); + let through = crate::topology::StepIdx(3); + let scheduler = RefillDrainScheduler::new( + Arc::clone(&refill), + through, + crate::topology::BranchIdx(0), + 100, + ); + assert_eq!(scheduler.walk(), WalkDirection::Reverse); + refill.store(true, Ordering::Relaxed); + assert_eq!(scheduler.walk(), WalkDirection::RefillThenReverse(through)); + refill.store(false, Ordering::Relaxed); + assert_eq!(scheduler.walk(), WalkDirection::Reverse); + } + + /// A byte-bounded queue stand-in whose depth the test sets directly. + struct FakeQueue(std::sync::atomic::AtomicU64); + + impl crate::queues::BoundedQueueHandle for FakeQueue { + fn current_bytes(&self) -> u64 { + self.0.load(std::sync::atomic::Ordering::Relaxed) + } + fn limit_bytes(&self) -> u64 { + u64::MAX + } + fn set_limit_bytes(&self, _new_limit: u64) {} + } + + fn registered( + producer: usize, + branch: usize, + queue: &std::sync::Arc, + ) -> crate::runtime::contexts::RegisteredQueue { + crate::runtime::contexts::RegisteredQueue { + producer_step_name: "producer", + producer_step: crate::topology::StepIdx(producer), + branch: crate::topology::BranchIdx(branch), + handle: std::sync::Arc::clone(queue) as _, + reorder_cap: None, + } + } + + /// Without a byte-bounded queue on the watched branch, `bind` leaves the + /// scheduler unbound: the read-ahead cap cannot be read, so the refill walk + /// follows the signal alone however full other queues are. `Debug` reports + /// the binding state. + #[test] + fn refill_drain_without_a_watched_queue_follows_the_signal_alone() { + use std::sync::Arc; + use std::sync::atomic::{AtomicBool, AtomicU64}; + let through = crate::topology::StepIdx(3); + let scheduler = RefillDrainScheduler::new( + Arc::new(AtomicBool::new(true)), + through, + crate::topology::BranchIdx(0), + 100, + ); + assert_eq!(scheduler.name(), "refill-drain"); + let full = Arc::new(FakeQueue(AtomicU64::new(1_000))); + scheduler.bind(&[registered(2, 0, &full), registered(3, 1, &full)]); + + assert!(scheduler.queue.get().is_none(), "no queue on the watched branch"); + assert_eq!(scheduler.walk(), WalkDirection::RefillThenReverse(through)); + let debug = format!("{scheduler:?}"); + assert!(debug.contains("bound: false"), "{debug}"); + + let ours = Arc::new(FakeQueue(AtomicU64::new(0))); + scheduler.bind(&[registered(3, 0, &ours)]); + let debug = format!("{scheduler:?}"); + assert!(debug.contains("bound: true"), "{debug}"); + } + + /// Re-binding the same queue is harmless (idempotent), but binding a + /// second run's different queue violates the one-run contract and trips + /// the debug assertion rather than silently watching the stale queue. + #[test] + fn refill_drain_rebind_to_the_same_queue_is_idempotent() { + use std::sync::Arc; + use std::sync::atomic::{AtomicBool, AtomicU64}; + let scheduler = RefillDrainScheduler::new( + Arc::new(AtomicBool::new(true)), + crate::topology::StepIdx(0), + crate::topology::BranchIdx(0), + 100, + ); + let ours = Arc::new(FakeQueue(AtomicU64::new(0))); + scheduler.bind(&[registered(0, 0, &ours)]); + scheduler.bind(&[registered(0, 0, &ours)]); + assert!(scheduler.queue.get().is_some(), "still bound after an idempotent rebind"); + } + + #[test] + #[cfg(debug_assertions)] + #[should_panic(expected = "bound to a second run's queue")] + fn refill_drain_rebind_to_another_queue_panics_in_debug() { + use std::sync::Arc; + use std::sync::atomic::{AtomicBool, AtomicU64}; + let scheduler = RefillDrainScheduler::new( + Arc::new(AtomicBool::new(true)), + crate::topology::StepIdx(0), + crate::topology::BranchIdx(0), + 100, + ); + let first = Arc::new(FakeQueue(AtomicU64::new(0))); + let second = Arc::new(FakeQueue(AtomicU64::new(0))); + scheduler.bind(&[registered(0, 0, &first)]); + scheduler.bind(&[registered(0, 0, &second)]); + } + + /// Once bound to the watched output branch, the refill walk stops at the + /// cap even while the signal is raised, and resumes below it. Queues from + /// another step, or from the refill step's other branch (a fan-out such as + /// a rejects output), are ignored. + #[test] + fn refill_drain_stops_reading_ahead_at_the_cap() { + use std::sync::Arc; + use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; + let through = crate::topology::StepIdx(3); + let scheduler = RefillDrainScheduler::new( + Arc::new(AtomicBool::new(true)), + through, + crate::topology::BranchIdx(1), + 100, + ); + let ours = Arc::new(FakeQueue(AtomicU64::new(0))); + let other = Arc::new(FakeQueue(AtomicU64::new(1_000))); + scheduler.bind(&[ + registered(2, 1, &other), + registered(3, 0, &other), + registered(3, 1, &ours), + ]); + + assert_eq!(scheduler.walk(), WalkDirection::RefillThenReverse(through)); + ours.0.store(100, Ordering::Relaxed); + assert_eq!(scheduler.walk(), WalkDirection::Reverse, "at the cap"); + ours.0.store(99, Ordering::Relaxed); + assert_eq!(scheduler.walk(), WalkDirection::RefillThenReverse(through)); + } } diff --git a/crates/fgumi-pipeline-core/src/step.rs b/crates/fgumi-pipeline-core/src/step.rs index 1e257a761..fb761adc2 100644 --- a/crates/fgumi-pipeline-core/src/step.rs +++ b/crates/fgumi-pipeline-core/src/step.rs @@ -184,6 +184,12 @@ pub struct StepProfile { #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum StepOutcome { /// Step did useful work (pushed an item or held one for later). + /// + /// Only a *new* hold counts: a retry of an already-held item that is still + /// rejected moved nothing, so report it as `NoProgress` (or `Contention`). + /// The round-robin walk restarts after a non-sticky step's `Progress`, so a + /// step that keeps reporting it while stuck is revisited forever and the + /// consumer that would drain its output is never reached. Progress, /// Step had nothing to do this call (input empty, no held work). NoProgress, diff --git a/src/lib/aligner.rs b/src/lib/aligner.rs index 6b4c1be46..14c8d14c7 100644 --- a/src/lib/aligner.rs +++ b/src/lib/aligner.rs @@ -16,8 +16,8 @@ //! [`substitute_template`] helper fills the `{ref}` / `{threads}` //! placeholders. //! -//! Used by the `AlignAndMergeStep` in -//! `src/lib/pipeline/steps/align_and_merge.rs`. This module is the +//! Used by the subprocess align backend's `SubprocessAlignStep` in +//! `src/lib/pipeline/steps/align/subprocess.rs`. This module is the //! framework-agnostic subprocess primitive; the typed `Step` impl that owns the //! I/O threads lives in that step module. @@ -702,23 +702,31 @@ impl Default for AlignerOptions { } } -/// Result of [`AlignerOptions::resolve`] — a ready-to-spawn aligner -/// invocation plus the parameters the downstream chain needs. +/// The alignment backend a [`ResolvedAligner`] will run. +#[derive(Debug, Clone)] +pub(crate) enum ResolvedBackend { + /// `--aligner::preset` or `--aligner::command`: the shell command to spawn + /// via [`AlignerProcess::spawn`] (already substituted and validated). + Subprocess { command: String }, +} + +/// Result of [`AlignerOptions::resolve`] — a ready-to-construct aligner +/// backend plus the parameters the downstream chain needs. /// /// `pub(crate)` because only the runall AAM dispatch consumes it; /// promote to `pub` if a cross-crate caller materializes. /// /// Every field is read when the chain builder wires the AAM stage -/// (`chains/builder.rs`): `command` seeds `AlignAndMergeStep`, `chunk_size` -/// derives the step's `in_flight_unmapped_budget` (via -/// `in_flight_budget_for_chunk_size`), and `mode`/`threads` are info-logged. +/// (`chains/builder.rs`): `backend` seeds the align step via +/// the align stage's `backend_for`, `chunk_size` derives the step's +/// `in_flight_unmapped_budget` (via `in_flight_budget_for_chunk_size`), and +/// `mode`/`threads` are info-logged. #[derive(Debug, Clone)] pub(crate) struct ResolvedAligner { - /// The shell command to spawn (already substituted and validated). - /// Passed directly to [`AlignerProcess::spawn`]. - pub command: String, + /// The backend to construct (by the align stage's `backend_for`). + pub backend: ResolvedBackend, /// The aligner's `-K` chunk size, in bases. The chain builder converts - /// it into `AlignAndMergeStep`'s `in_flight_unmapped_budget` so the + /// it into the align backend's in-flight byte budget so the /// resident unmapped backlog stays proportional to one aligner chunk. pub chunk_size: u64, /// Resolved thread count (default-substituted for preset mode; @@ -745,7 +753,7 @@ pub(crate) enum ResolvedAlignerMode { impl AlignerOptions { /// Validate the option combination and produce a [`ResolvedAligner`] - /// ready for [`AlignerProcess::spawn`]. + /// ready for the align stage's `backend_for` to construct. /// /// # Arguments /// @@ -803,7 +811,7 @@ impl AlignerOptions { let command = preset.build_command(reference, threads, self.chunk_size, aligner_bin); Ok(ResolvedAligner { - command, + backend: ResolvedBackend::Subprocess { command }, chunk_size: self.chunk_size, threads: Some(threads), mode: ResolvedAlignerMode::Preset(preset), @@ -838,7 +846,7 @@ impl AlignerOptions { // hardcode any thread count). let command = substitute_template(&template, reference, top_threads)?; Ok(ResolvedAligner { - command, + backend: ResolvedBackend::Subprocess { command }, chunk_size: self.chunk_size, threads: None, mode: ResolvedAlignerMode::Command, @@ -1005,7 +1013,7 @@ mod tests { /// `kill()` reports `true` when the child is reaped within the deadline. /// A plain `sleep` responds to SIGKILL immediately, so the bounded kill /// reaps it well inside the 1s window — and a `true` return is what lets - /// `AlignAndMergeStep::Drop` safely issue its follow-up `wait()` without + /// `SubprocessAlignStep::drop` safely issue its follow-up `wait()` without /// risking a teardown hang. #[test] fn test_kill_reports_reaped_for_killable_child() { @@ -1141,8 +1149,9 @@ mod tests { chunk_size: DEFAULT_ALIGNER_CHUNK_SIZE, }; let resolved = opts.resolve(&ref_path, 4, None).unwrap(); - assert!(resolved.command.contains(&ref_path.display().to_string())); - assert!(resolved.command.contains("-t 4")); + let ResolvedBackend::Subprocess { command } = resolved.backend; + assert!(command.contains(&ref_path.display().to_string())); + assert!(command.contains("-t 4")); assert!(matches!(resolved.mode, ResolvedAlignerMode::Command)); } diff --git a/src/lib/commands/zipper.rs b/src/lib/commands/zipper.rs index b54a0ec97..20a15f699 100644 --- a/src/lib/commands/zipper.rs +++ b/src/lib/commands/zipper.rs @@ -527,7 +527,7 @@ fn collect_mapped_indices(mapped: &Template, is_first_segment: bool) -> Vec { /// `HeaderHandle` stashed by [`Self::add_align`] for consumption by /// [`Self::add_sink`]. /// - /// `AlignAndMergeStep` resolves the final merged header at runtime - /// (after the aligner emits its own `@PG`/`@RG`/`@CO` lines). The - /// downstream [`WriteBgzfFile`] must therefore be constructed with + /// The align backend resolves the final merged header at runtime (after + /// the aligner emits its own `@PG`/`@RG`/`@CO` lines). The downstream + /// [`WriteBgzfFile`] must therefore be constructed with /// [`WriteBgzfFile::new_with_handle`], which blocks until the handle /// is resolved before writing the BAM header bytes. /// @@ -2416,7 +2416,7 @@ impl<'a> ChainBuilder<'a> { /// /// For [`StagePosition::Intermediate`], `SerializeBamRecords` is **not** /// appended: the chain tail is left as `BamTemplateBatch` so the next stage - /// (`add_align` → `AlignAndMergeStep`) can consume the correct step's kept + /// (`add_align` → the align backend) can consume the correct step's kept /// output (branch 0) directly; `add_align` skips `GroupByQueryname` on an /// incoming `BamTemplateBatch`. This is the correct→align fused path. /// @@ -2478,7 +2478,7 @@ impl<'a> ChainBuilder<'a> { // (the `--start-from correct --stop-after ≥ zipper` fused chain). // In that case we skip `SerializeBamRecords` and leave the tail at // the correct step's `BamTemplateBatch` output (branch 0) so - // `add_align` can wire `GroupByQueryname → AlignAndMergeStep` next. + // `add_align` can wire the align backend directly next. let correct_opts = self .spec @@ -2612,8 +2612,8 @@ impl<'a> ChainBuilder<'a> { self.chain_tail_kind = ChainTailKind::SerializedBytes; } else { // Intermediate: leave the tail as BamTemplateBatch so the next - // stage (add_align → GroupByQueryname → AlignAndMergeStep) can - // consume it directly. + // stage (add_align → the align backend, skipping GroupByQueryname) + // can consume it directly. self.current_tail = Some(process_tail); self.chain_tail_kind = ChainTailKind::BamTemplateBatch; } @@ -2644,15 +2644,16 @@ impl<'a> ChainBuilder<'a> { Ok(()) } - /// AlignAndMerge-specific step sequence (subprocess source): + /// AlignAndMerge-specific step sequence: /// /// ```text /// (source preamble: ReadBgzfBlocks → BgzfDecompress → FindBamBoundaries → DecodeRecords) /// ↓ - /// GroupByQueryname ← appended here - /// ↓ - /// AlignAndMergeStep ← appended here (spawns aligner subprocess) + /// GroupByQueryname ← appended here (skipped when the upstream tail is BamTemplateBatch) /// ↓ + /// ← appended here: `SubprocessAlignStep` + /// ↓ (spawns the aligner subprocess) followed by + /// ↓ the shared `MergeAlignedStep` /// (Terminal) SerializeBamRecords ← appended only when Terminal /// ↓ /// (next stage consumes BamTemplateBatch — Intermediate only) @@ -2661,8 +2662,8 @@ impl<'a> ChainBuilder<'a> { /// Both [`StagePosition::Terminal`] and [`StagePosition::Intermediate`] /// are supported: /// - /// - **Terminal** — `SerializeBamRecords` is appended after the AAM step - /// so the chain tail is `DecompressedBlock` ready for + /// - **Terminal** — `SerializeBamRecords` is appended after the merged + /// align output so the chain tail is `DecompressedBlock` ready for /// `BgzfCompress → WriteBgzfFile`. This is the `--stop-after zipper` /// path when `--start-from align-and-merge`: the merged BAM is /// written directly to the output file. @@ -2672,16 +2673,16 @@ impl<'a> ChainBuilder<'a> { /// ## Header handling /// /// The aligner contributes `@PG`/`@RG`/`@CO` lines to the output - /// BAM header at runtime (not at chain-construction time). The - /// `AlignAndMergeStep` resolves a [`HeaderHandle`] once the aligner - /// emits its SAM/BAM header; the downstream [`WriteBgzfFile`] - /// (added by [`Self::add_sink`]) blocks on the handle before - /// writing any record bytes. + /// BAM header at runtime (not at chain-construction time). The align + /// backend resolves a [`HeaderHandle`] with the merged header once the + /// aligner emits its SAM/BAM header. The downstream + /// [`WriteBgzfFile`] (added by [`Self::add_sink`]) blocks on the handle + /// before writing any record bytes. /// /// To prepare the correct partial header (dict-derived `@SQ` + - /// unmapped-BAM headers + fgumi `@PG`) for `AlignAndMergeConfig`, - /// `add_align` calls [`build_output_header`] with the current - /// `self.header` (which was built from the unmapped BAM by + /// unmapped-BAM headers + fgumi `@PG`) for the backend and the shared + /// `MergeConfig`, `add_align` calls [`build_output_header`] with the + /// current `self.header` (which was built from the unmapped BAM by /// `ChainBuilder::new`) and the reference FASTA's `.dict` path. /// The result replaces `self.header` so downstream stages and /// `add_sink` see the dict-derived `@SQ` table. @@ -2691,24 +2692,18 @@ impl<'a> ChainBuilder<'a> { /// /// ## Thread floor /// - /// AAM spawns two daemon threads (FASTQ writer + SAM/BAM reader) plus - /// the aligner subprocess. The pipeline needs at least 4 framework - /// workers for steady-state progress: - /// - /// - 1 worker to drive the source preamble - /// - 1 worker dispatching `AlignAndMergeStep` (Serial) - /// - 1+ workers for any downstream Parallel steps - /// - 1 spare for progress / bookkeeping - /// - /// `override_pipeline_threads` is raised to `threads.max(4)` so - /// `build()` forwards the correct value to `PipelineConfig::threads`. + /// The backend reports its worker floor in `AlignWired::min_workers`, and + /// `fold_align_wired_scheduling` raises `override_pipeline_threads` to + /// `num_threads.max(min_workers)` so `build()` forwards the correct value + /// to `PipelineConfig::threads`. The subprocess backend needs 4 (it spawns + /// two daemon threads plus the aligner subprocess). /// /// ## Errors /// /// Returns errors if: /// - align options are missing from the spec bag, /// - the reference `.dict` file cannot be found, - /// - `AlignAndMergeStep::new` fails to spawn the aligner subprocess. + /// - `SubprocessAlignStep::new` fails to spawn the aligner subprocess. /// /// [`HeaderHandle`]: crate::pipeline::core::header::HeaderHandle /// [`WriteBgzfFile`]: crate::pipeline::steps::sink::write_bgzf::WriteBgzfFile @@ -2720,7 +2715,8 @@ impl<'a> ChainBuilder<'a> { use crate::logging::OperationTimer; use crate::pipeline::chains::commands::align::AlignFinalizeHook; use crate::pipeline::core::header::HeaderHandle; - use crate::pipeline::steps::align_and_merge::{AlignAndMergeConfig, AlignAndMergeStep}; + use crate::pipeline::steps::align::merge::MergeConfig; + use crate::pipeline::steps::align::{AlignBackend, AlignWiringCtx, backend_for}; use crate::pipeline::steps::serialize::SerializeBamRecords; use crate::reference::find_dict_path; use crate::sam::check_sort; @@ -2748,7 +2744,6 @@ impl<'a> ChainBuilder<'a> { "AlignAndMerge: mode = {:?}, chunk_size = {}, threads = {:?}", resolved.mode, resolved.chunk_size, resolved.threads, ); - info!("AlignAndMerge command: {}", resolved.command); // Validate the unmapped BAM source has queryname sort order — but only // when Align reads directly from the chain source (a file). When an @@ -2790,45 +2785,50 @@ impl<'a> ChainBuilder<'a> { let header_handle = HeaderHandle::new(); let records_emitted = std::sync::Arc::new(AtomicU64::new(0)); - let aam_cfg = AlignAndMergeConfig { - tag_info: std::sync::Arc::new(TagInfo::new(Vec::new(), Vec::new(), Vec::new())), - skip_tc_tags: false, - reference: None, - partial_output_header: std::sync::Arc::clone(&partial_header), - header_handle: header_handle.clone(), - records_emitted: std::sync::Arc::clone(&records_emitted), - output_byte_limit: self.tuning.per_step_byte_limit, - // Bound AAM's in-flight unmapped backlog to the aligner's `-K` - // chunk size so a fast-draining aligner can't accumulate the - // whole input's unmapped reads in RAM (issue #382). - in_flight_unmapped_budget: - crate::pipeline::steps::align_and_merge::in_flight_budget_for_chunk_size( - resolved.chunk_size, - ), - }; - - let aam_step = AlignAndMergeStep::new(aam_cfg, &resolved.command) - .map_err(|e| anyhow!("AlignAndMergeStep::new: {e:#}"))?; - let tail = self.current_tail.expect("add_align called before add_source"); - // Wire GroupByQueryname before the AAM step — but only when the upstream - // stage emits DecodedRecordBatch (the normal source preamble path). - // When `chain_tail_kind == BamTemplateBatch`, an upstream stage (e.g. - // Correct) has already grouped records into BamTemplateBatch, which is - // exactly the input type AlignAndMergeStep expects. Skip GroupByQueryname - // in that case to avoid a DecodedRecordBatch→BamTemplateBatch type mismatch. + // Wire GroupByQueryname before the align backend — but only when the + // upstream stage emits DecodedRecordBatch (the normal source preamble + // path). When `chain_tail_kind == BamTemplateBatch`, an upstream stage + // (e.g. Correct) has already grouped records into BamTemplateBatch, which + // is exactly the input type the align backend expects. Skip + // GroupByQueryname in that case to avoid a + // DecodedRecordBatch→BamTemplateBatch type mismatch. let align_input_tail = if self.chain_tail_kind == ChainTailKind::BamTemplateBatch { - // Upstream is already BamTemplateBatch — wire AAM directly. + // Upstream is already BamTemplateBatch — wire the backend directly. tail } else { // Upstream is DecodedRecordBatch (source preamble). Group by queryname - // to produce BamTemplateBatch for AlignAndMergeStep — the parallel + // to produce BamTemplateBatch for the align backend — the parallel // `AssembleTemplates` map when the source cut left batches closed // under queryname, else the serial `GroupByQueryname`. self.append_queryname_grouper(tail) }; - let aam_tail = self.pipeline.append_step(aam_step, align_input_tail); + + // Construct the align backend `resolve` selected behind the + // `AlignBackend` trait, wire its steps after the grouped input, and + // obtain the merged `BamTemplateBatch` tail plus the scheduling facts to + // fold in. + let backend: Box = backend_for(resolved); + info!("AlignAndMerge backend: {}", backend.describe()); + // The zipper merge every aligned template goes through, run by the + // backend (the subprocess backend's shared `Parallel` merge step). + let merge_config = std::sync::Arc::new(MergeConfig { + tag_info: std::sync::Arc::new(TagInfo::new(Vec::new(), Vec::new(), Vec::new())), + skip_tc_tags: false, + reference: None, + partial_output_header: std::sync::Arc::clone(&partial_header), + records_emitted: std::sync::Arc::clone(&records_emitted), + output_byte_limit: self.tuning.per_step_byte_limit, + }); + let wiring_ctx = AlignWiringCtx { + partial_output_header: std::sync::Arc::clone(&partial_header), + header_handle: header_handle.clone(), + per_step_byte_limit: self.tuning.per_step_byte_limit, + merge: merge_config, + }; + let wired = backend.wire(&self.pipeline, align_input_tail, &wiring_ctx)?; + let aam_tail = wired.tail; if position == StagePosition::Terminal { // Terminal: append SerializeBamRecords so the chain tail is @@ -2854,14 +2854,21 @@ impl<'a> ChainBuilder<'a> { self.pending_header_transform = None; self.pending_header_handle = Some(header_handle); - // AAM spawns two daemon threads + aligner subprocess; 4 workers - // minimum for steady-state progress. - let effective = num_threads.max(4); - if let Some(prev) = self.override_pipeline_threads { - self.override_pipeline_threads = Some(prev.max(effective)); - } else { - self.override_pipeline_threads = Some(effective); - } + // Raise the pool floor to the backend's `min_workers` (the subprocess + // backend needs 4: it spawns two daemon threads + the aligner subprocess, + // so 4 framework workers are the floor for steady-state progress) — and + // fold its scheduler preference into the automatic decision. Only ever set the flag true (mirrors the + // grouping stages); the subprocess backend leaves it false, preserving + // today's behavior. See `fold_align_wired_scheduling`'s doc comment for + // the full rationale (it is unit-tested there without needing a real + // backend). + fold_align_wired_scheduling( + num_threads, + wired.min_workers, + wired.prefers_drain_first, + &mut self.override_pipeline_threads, + &mut self.use_drain_first_scheduler, + ); let timer = OperationTimer::new("AlignAndMerge"); // Success-only: the hook logs "pipeline completed successfully" and @@ -6187,6 +6194,40 @@ fn grouping_stage_wants_drain_first(stage: Stage, position: StagePosition) -> bo ) } +/// Folds an align backend's wiring result (`AlignWired::min_workers`, +/// `AlignWired::prefers_drain_first`) into the chain-wide +/// thread floor and the accumulated drain-first flag. +/// +/// `num_threads.max(wired_min_workers)` raises the floor to whichever is +/// larger, and never lowers a floor a previous stage already raised (mirrors +/// `add_sort`/`add_zipper`'s existing "max, don't overwrite" thread-override +/// discipline). `wired_prefers_drain_first` only ever sets the flag `true`, +/// never clears it — a previous stage's opt-in survives a later stage that +/// itself does not need drain-first, matching `grouping_stage_wants_drain_first`. +/// +/// The subprocess backend's floor of 4 forces `n_threads >= 4` for any align +/// chain, which forecloses `crate::pipeline::core::runtime::fused:: +/// should_fuse_single_thread`'s `n_threads == 1` precondition regardless of the +/// user's `--threads` request; a backend with a floor of 1 would not. +/// +/// Extracted as a free function, taking primitives rather than `&AlignWired`, +/// so the fold is unit-testable without constructing a real align backend +/// (`AlignBackend::wire()` needs a spawnable aligner subprocess). +fn fold_align_wired_scheduling( + num_threads: usize, + wired_min_workers: usize, + wired_prefers_drain_first: bool, + override_pipeline_threads: &mut Option, + use_drain_first_scheduler: &mut bool, +) { + let effective = num_threads.max(wired_min_workers); + *override_pipeline_threads = + Some(override_pipeline_threads.map_or(effective, |prev| prev.max(effective))); + if wired_prefers_drain_first { + *use_drain_first_scheduler = true; + } +} + /// Resolve the effective drain-first choice from the hidden `--pool-scheduler` /// override and the automatic per-chain decision. /// @@ -6253,6 +6294,7 @@ fn dict_not_found_error(reference: &std::path::Path) -> anyhow::Error { mod tests { use super::*; use crate::commands::common::PoolScheduler; + use crate::pipeline::steps::align::subprocess::SubprocessBackend; /// A `ChainSpec` with the given stages and every other field at its /// simplest valid default. Mirrors `validate::tests::empty_spec` — kept as @@ -6515,6 +6557,90 @@ mod tests { ); } + /// Pin `fold_align_wired_scheduling`. The `subprocess_*` cases read + /// `SubprocessBackend::MIN_WORKERS` / `PREFERS_DRAIN_FIRST` directly, so a + /// change to those constants is seen here; the `floor_1_*` cases cover a + /// backend reporting a floor of 1 that prefers drain-first, which leaves the + /// resolved floor at the requested thread count. None needs a real aligner + /// subprocess (`AlignBackend::wire()` does, so it cannot run in this unit + /// test). + #[rstest::rstest] + #[case::floor_1_at_requested_threads_1_keeps_the_floor_at_1( + 1, + 1, + true, + None, + false, + Some(1), + true + )] + #[case::floor_1_at_requested_threads_8_raises_the_floor_to_8( + 8, + 1, + true, + None, + false, + Some(8), + true + )] + #[case::subprocess_at_requested_threads_1_forces_the_floor_to_4( + 1, + SubprocessBackend::MIN_WORKERS, + SubprocessBackend::PREFERS_DRAIN_FIRST, + None, + false, + Some(4), + false + )] + #[case::subprocess_never_lowers_a_higher_existing_floor( + 1, + SubprocessBackend::MIN_WORKERS, + SubprocessBackend::PREFERS_DRAIN_FIRST, + Some(6), + false, + Some(6), + false + )] + #[case::higher_min_workers_raises_a_lower_existing_floor( + 1, + 8, + true, + Some(4), + false, + Some(8), + true + )] + #[case::drain_first_is_sticky_once_set_by_an_earlier_stage( + 4, + SubprocessBackend::MIN_WORKERS, + SubprocessBackend::PREFERS_DRAIN_FIRST, + None, + true, + Some(4), + true + )] + fn fold_align_wired_scheduling_matrix( + #[case] num_threads: usize, + #[case] wired_min_workers: usize, + #[case] wired_prefers_drain_first: bool, + #[case] initial_override: Option, + #[case] initial_drain_first: bool, + #[case] expected_override: Option, + #[case] expected_drain_first: bool, + ) { + let mut override_pipeline_threads = initial_override; + let mut use_drain_first_scheduler = initial_drain_first; + fold_align_wired_scheduling( + num_threads, + wired_min_workers, + wired_prefers_drain_first, + &mut override_pipeline_threads, + &mut use_drain_first_scheduler, + ); + assert_eq!(override_pipeline_threads, expected_override, "thread floor mismatch"); + assert_eq!(use_drain_first_scheduler, expected_drain_first, "drain-first flag mismatch"); + } + /// Pin the hidden `--pool-scheduler` override resolution: `Auto` defers to /// the automatic per-chain decision, while `DrainFirst`/`ChainOrder` force /// their direction regardless of it. This is the A/B knob's whole contract, diff --git a/src/lib/pipeline/chains/commands/align.rs b/src/lib/pipeline/chains/commands/align.rs index a247cc5a6..d14a0b54a 100644 --- a/src/lib/pipeline/chains/commands/align.rs +++ b/src/lib/pipeline/chains/commands/align.rs @@ -23,8 +23,8 @@ use crate::pipeline::chains::FinalizeHook; /// `timer.log_completion` so the per-stage wall-time line appears in /// the log. pub(crate) struct AlignFinalizeHook { - /// Counter atomically incremented by `AlignAndMergeStep` as it - /// emits `BamTemplateBatch`es downstream. + /// Counter atomically incremented by `merge::merge_zipper_batch` as + /// merged `BamTemplateBatch`es go downstream. pub(crate) records_emitted: Arc, pub(crate) timer: OperationTimer, } diff --git a/src/lib/pipeline/chains/stage.rs b/src/lib/pipeline/chains/stage.rs index 70c7190d6..79b123377 100644 --- a/src/lib/pipeline/chains/stage.rs +++ b/src/lib/pipeline/chains/stage.rs @@ -50,7 +50,7 @@ impl Stage { } /// Returns `true` if this stage produces `BamTemplateBatch` from - /// an external source (AAM subprocess or zipper merge). + /// an external source (the align stage's aligner or zipper merge). #[must_use] pub fn is_template_producer(self) -> bool { matches!(self, Stage::Align | Stage::Zipper) diff --git a/src/lib/pipeline/steps/align/merge.rs b/src/lib/pipeline/steps/align/merge.rs new file mode 100644 index 000000000..33029d60b --- /dev/null +++ b/src/lib/pipeline/steps/align/merge.rs @@ -0,0 +1,400 @@ +//! `MergeAlignedStep` — the shared, `Parallel` consumer of an align backend's +//! `ZipperBatch` stream. +//! +//! Each `ZipperBatch` carries both halves of a batch (aligner-emitted `mapped` +//! templates positionally paired with the original `unmapped` templates) plus a +//! dense serial, so the merge is per-template and needs no cross-batch state: +//! `merge_one_template_with` (`merge_raw_with` plus optional bisulfite restore) +//! is run for each pair, folding record-count and heap-size accounting into the +//! same pass. That independence is why this step is `Parallel` with an +//! `OrderedBytesSingle` + `BranchOrdering::ByItemOrdinal` +//! output: workers merge batches concurrently and the framework restores input +//! order from each `BamTemplateBatch`'s serial. +//! +//! This was previously the `merge_zipper_batch` call inlined on the (Serial) +//! `AlignAndMergeStep`'s dispatching worker; promoting it to its own `Parallel` +//! step is the split the module doc on the former `align_and_merge.rs` +//! anticipated. The merge logic itself is unchanged. + +use std::io; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; + +use noodles::sam::Header; + +use crate::commands::zipper::{ZipperTags, merge_one_template_with}; +use crate::pipeline::core::Unpushed; +use crate::pipeline::core::held::HeldSlot; +use crate::pipeline::core::outputs::OrderedBytesSingle; +use crate::pipeline::core::queues::QueueSpec; +use crate::pipeline::core::reorder::BranchOrdering; +use crate::pipeline::core::step::{Step, StepCtx, StepKind, StepOutcome, StepProfile}; +use crate::pipeline::steps::align::ZipperBatch; +use crate::pipeline::steps::types::BamTemplateBatch; +use crate::reference::ReferenceReader; +use crate::template::Template; +use crate::umi::TagInfo; + +/// Immutable, Arc-friendly configuration for [`MergeAlignedStep`]. +/// +/// One `Arc` clone is shared by every per-worker copy of the step; the cost is +/// negligible compared to the data flowing through. +pub(crate) struct MergeConfig { + /// Tag-merge rules (remove / reverse / revcomp) plumbed to `merge_raw`. + pub(crate) tag_info: Arc, + + /// Whether to skip TC (template-coordinate) tag handling. Mirrors + /// `ZipperMergeConfig::skip_tc_tags`. + pub(crate) skip_tc_tags: bool, + + /// Optional reference reader used by + /// `restore_unconverted_bases_in_raw_template` (bisulfite path). + /// `None` for normal alignment. + pub(crate) reference: Option>, + + /// Partial output header (dict-derived `@SQ` + unmapped-derived + /// `@HD`/`@CO`/`@RG`/`@PG` + fgumi `@PG`) used by `merge_one_template_with`. + pub(crate) partial_output_header: Arc
, + + /// Counter for records emitted downstream. Exposed back to the caller after + /// `Pipeline::run` returns so summary logging can report a real throughput + /// number. + pub(crate) records_emitted: Arc, + + /// Byte limit for the downstream output queue (`OrderedBytesSingle` + /// `ByteBounded`) that merged `BamTemplateBatch`es are pushed onto. + pub(crate) output_byte_limit: u64, +} + +/// `Parallel + ByItemOrdinal` step that merges [`ZipperBatch`]es into +/// `BamTemplateBatch`es via [`merge_zipper_batch`]. +pub(crate) struct MergeAlignedStep { + cfg: Arc, + /// Precomputed tag-merge bitsets, built once from `cfg.tag_info` and reused + /// for every template across every batch — `cfg.tag_info` is immutable for + /// the step's whole lifetime, so there is no reason to rebuild `ZipperTags` + /// (three `TagBitset` allocations) per batch, let alone per template. + tags: Arc, + held: HeldSlot>, +} + +impl MergeAlignedStep { + /// Build the shared merge step from `cfg`, precomputing the `ZipperTags` + /// bitsets once. + pub(crate) fn from_shared(cfg: Arc) -> Self { + let tags = Arc::new(ZipperTags::from_tag_info(&cfg.tag_info)); + Self { cfg, tags, held: HeldSlot::new() } + } +} + +impl Step for MergeAlignedStep { + type Input = ZipperBatch; + type Outputs = OrderedBytesSingle; + + fn profile(&self) -> StepProfile { + StepProfile { + name: "MergeAligned", + kind: StepKind::Parallel, + sticky: false, + output_queues: vec![QueueSpec::ByteBounded { limit_bytes: self.cfg.output_byte_limit }], + branch_ordering: vec![BranchOrdering::ByItemOrdinal], + } + } + + fn try_run(&mut self, ctx: &mut StepCtx<'_, Self>) -> io::Result { + // Drain the held output slot first, before popping more input. + if let Some(unpushed) = self.held.take() { + match ctx.outputs.retry(unpushed) { + Ok(()) => {} + Err(again) => { + self.held.put(again); + // `Contention` (not `NoProgress`) — returning `NoProgress` + // from a Parallel step on a held slot risks the framework + // marking this worker `Skip` if upstream is also drained, + // silently dropping the held item. Mirrors + // `TemplatesToRecordBatch` / `BgzfDecompress`. + return Ok(StepOutcome::Contention); + } + } + } + + let Some(zb) = ctx.input.pop() else { + // No input this call. If upstream is drained, every item has been + // processed (held output flushed above) and this step will never + // push again — report Finished. For a Parallel step only the last + // clone to finish closes the shared output (gated by the driver's + // StepDrainCounter). + if ctx.input.is_drained() { + return Ok(StepOutcome::Finished); + } + return Ok(StepOutcome::NoProgress); + }; + + let merged = merge_zipper_batch(zb, &self.cfg, &self.tags)?; + match ctx.outputs.push(merged) { + Ok(()) => Ok(StepOutcome::Progress), + Err(unpushed) => { + self.held.put(unpushed); + Ok(StepOutcome::Progress) + } + } + } + + fn new_worker_copy(&self) -> Self { + Self { cfg: Arc::clone(&self.cfg), tags: Arc::clone(&self.tags), held: HeldSlot::new() } + } +} + +/// Merge one `ZipperBatch` into a `BamTemplateBatch`. Takes the caller's +/// precomputed `tags` (built once for the whole step, since `cfg.tag_info` is +/// immutable) and, per template, runs `merge_one_template_with` +/// (`merge_raw_with` plus optional bisulfite restore); folds record-count and +/// heap-size accounting into the same single pass so the resulting +/// `BamTemplateBatch` doesn't re-walk the templates to compute `total_bytes`. +pub(crate) fn merge_zipper_batch( + zb: ZipperBatch, + cfg: &MergeConfig, + tags: &ZipperTags, +) -> io::Result { + let ZipperBatch { serial, mapped, unmapped } = zb; + // Release-safe guard (not a debug_assert): a violated length invariant would + // otherwise make the `zip` below silently truncate to the shorter side, + // dropping templates with no error. Fail loudly instead — mirroring + // `extract_batch`'s record-count guard. (A bare debug_assert would also + // shadow this Err path from CI's debug-build tests.) + if mapped.len() != unmapped.templates().len() { + return Err(io::Error::other(format!( + "align-and-merge: ZipperBatch invariant violated — {} mapped vs {} unmapped \ + templates; refusing to zip mismatched halves", + mapped.len(), + unmapped.templates().len(), + ))); + } + + let mut merged: Vec