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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions crates/fgumi-pipeline-core/src/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion crates/fgumi-pipeline-core/src/header.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//!
Expand Down
219 changes: 203 additions & 16 deletions crates/fgumi-pipeline-core/src/runtime/driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -405,6 +426,7 @@ fn round_robin_dispatch(
// safe and avoids reorder-under-iteration.
let mut finished: Vec<StepIdx> = Vec::new();
let n = live.len();
let split = refill_split(walk, live.order());
for i in 0..n {
if signal.is_done() {
break;
Expand All @@ -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;
Expand All @@ -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
Expand Down Expand Up @@ -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<StepIdx> = [0, 2, 3, 5, 7].into_iter().map(StepIdx).collect();
let split = refill_split(walk, &order);
let visited: Vec<usize> =
(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;
Expand Down Expand Up @@ -1398,18 +1450,20 @@ 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
// ever reached). Removing either one hangs `sticky_holding_source_yields_to_
// 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<u32>,
Expand Down Expand Up @@ -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<u32>,
}
impl Step for HoldingSource {
type Input = ();
type Outputs = Single<u32>;
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<StepOutcome> {
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<Box<dyn ErasedStep>> = 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<WorkerStepEntry> =
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::<Vec<u32>>(),
"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)]
Expand Down
4 changes: 3 additions & 1 deletion crates/fgumi-pipeline-core/src/runtime/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Loading
Loading