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
3 changes: 2 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -165,7 +165,8 @@ because the hot path runs once per record (BAM workloads are millions to billion
of records) and a safe rewrite measurably regresses sort throughput.

- **`crates/fgumi-sort/src/inline.rs`** — three `#[allow(unsafe_code)]` regions:
the `radix_sort_record_refs` (coordinate) and `radix_sort_template_refs` /
the `radix_sort_record_refs_with_max` (coordinate; `radix_sort_record_refs`
delegates to it after scanning for the max) and `radix_sort_template_refs` /
`radix_sort_template_field` (template) LSD radix sorts use `Vec::set_len` to
skip per-element initialization on the auxiliary scratch buffer, plus
raw-pointer slice swaps to avoid double-borrow restrictions across the
Expand Down
101 changes: 101 additions & 0 deletions benches/core_functions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1305,6 +1305,106 @@ fn bench_queryname_sort_strategies(c: &mut Criterion) {
group.finish();
}

/// Coordinate-sort radix throughput across the three ways `bytes_needed` can be
/// sized, isolating the two independent wins:
///
/// - `full_width` reproduces the pre-optimization behavior by passing the
/// unmapped sentinel as the bound, forcing all 8 radix passes — which is what
/// a sentinel-inclusive max scan yields the moment a single unmapped read is
/// present, i.e. on essentially every real coordinate-sorted BAM.
/// - `scan_mapped_max` scans for the largest *mapped* key, sizing the passes to
/// the ~4–5 bytes such a key occupies. The delta against `full_width` is the
/// pass-count win.
/// - `tracked_max` supplies that same bound without scanning, using a value the
/// benchmark precomputed while building the input. The delta against
/// `scan_mapped_max` is the scan-elimination win, and it is small enough that
/// `RecordBuffer` deliberately does not chase it: it scans once per sort
/// rather than maintaining a running bound across pushes.
///
/// `sort_unstable` is a comparison-sort baseline. Note it is *unstable*, so it
/// is not a drop-in substitute — coordinate sort order must be stable to match
/// `samtools sort`. The input mixes 25 tids with a 5% unmapped tail to exercise
/// the sentinel handling.
fn bench_coordinate_radix_sort(c: &mut Criterion) {
use fgumi_sort::{
PackedCoordinateKey, RecordRef, radix_sort_record_refs, radix_sort_record_refs_with_max,
};

let nref = 25u32; // chr1-22, X, Y, MT
let mut group = c.benchmark_group("coordinate_radix_sort");

for &n in &[1_000_000usize, 8_000_000] {
// Deterministic xorshift so the dataset is reproducible without an RNG dep.
let mut state = 0x9e37_79b9_7f4a_7c15u64;
let mut next = move || {
state ^= state << 13;
state ^= state >> 7;
state ^= state << 17;
state
};

let mut refs: Vec<RecordRef> = Vec::with_capacity(n);
let mut mapped_max = 0u64;
for _ in 0..n {
let r = next();
let key = if r % 20 == 0 {
u64::MAX // ~5% unmapped (sentinel)
} else {
let tid = i32::try_from(r % u64::from(nref)).unwrap();
let pos = i32::try_from((r >> 16) % 250_000_000).unwrap();
let k = PackedCoordinateKey::new(tid, pos, false, nref).0;
mapped_max = mapped_max.max(k);
k
};
refs.push(RecordRef::new(key, 0, 0));
}

group.throughput(Throughput::Elements(n as u64));
group.bench_with_input(BenchmarkId::new("full_width", n), &refs, |b, refs| {
b.iter_batched(
|| refs.clone(),
|mut v| {
radix_sort_record_refs_with_max(&mut v, u64::MAX);
black_box(v)
},
criterion::BatchSize::LargeInput,
);
});
group.bench_with_input(BenchmarkId::new("scan_mapped_max", n), &refs, |b, refs| {
b.iter_batched(
|| refs.clone(),
|mut v| {
radix_sort_record_refs(&mut v);
black_box(v)
},
criterion::BatchSize::LargeInput,
);
});
group.bench_with_input(BenchmarkId::new("tracked_max", n), &refs, |b, refs| {
b.iter_batched(
|| refs.clone(),
|mut v| {
radix_sort_record_refs_with_max(&mut v, mapped_max);
black_box(v)
},
criterion::BatchSize::LargeInput,
);
});
group.bench_with_input(BenchmarkId::new("sort_unstable", n), &refs, |b, refs| {
b.iter_batched(
|| refs.clone(),
|mut v| {
v.sort_unstable_by_key(|r| r.sort_key);
black_box(v)
},
criterion::BatchSize::LargeInput,
);
});
}

group.finish();
}

criterion_group!(
benches,
bench_phred_conversions,
Expand All @@ -1318,5 +1418,6 @@ criterion_group!(
bench_vanilla_consensus_caller,
bench_queryname_comparators,
bench_queryname_sort_strategies,
bench_coordinate_radix_sort,
);
criterion_main!(benches);
3 changes: 2 additions & 1 deletion crates/fgumi-raw-bam/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -98,7 +98,8 @@ pub use overlap::{

// -- raw_bam_record --
pub use raw_bam_record::{
RawBamReader, RawBamWriter, RawRecord, read_raw_record, write_raw_record, write_raw_records,
RawBamReader, RawBamWriter, RawRecord, read_block_size, read_raw_record, write_raw_record,
write_raw_records,
};

// -- sequence --
Expand Down
18 changes: 15 additions & 3 deletions crates/fgumi-raw-bam/src/raw_bam_record.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,10 +171,22 @@ where
Ok(block_size)
}

/// Reads the 4-byte block size prefix.
/// Reads the 4-byte BAM record `block_size` prefix (little-endian u32).
///
/// Returns 0 at EOF (no bytes available).
fn read_block_size<R>(reader: &mut R) -> io::Result<usize>
/// Returns 0 at EOF (no bytes available before the prefix starts). Handles a
/// prefix split across reader boundaries by issuing a 1-byte read first (to
/// detect clean EOF) followed by a `read_exact` of the remaining 3 bytes.
///
/// Exposed for the sort ingest borrow-in-place path
/// (`PooledInputStream::next_record_borrowed`), which reads the prefix itself so
/// it can decide whether the record body can be borrowed from the current
/// decompressed block or must be gathered into a scratch buffer.
///
/// # Errors
///
/// Returns an error if the reader errors, the prefix is truncated, or the value
/// does not fit in `usize`.
pub fn read_block_size<R>(reader: &mut R) -> io::Result<usize>
where
R: Read,
{
Expand Down
14 changes: 10 additions & 4 deletions crates/fgumi-sort/src/external.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2263,12 +2263,16 @@ impl RawExternalSorter {
debug!("Phase 1: Reading and sorting chunks (inline buffer, keyed output)...");
let mut probe = SpillProbe::new("phase1");

for record in record_source.by_ref() {
// Borrow each record's bytes straight out of the decompressed block and
// push them into the arena, skipping the intermediate `RawRecord` copy
// (the borrowed slice is invalidated by the next `next_record_borrowed`
// call, which is fine — `push_coordinate` copies the bytes into the buffer).
while let Some(record) = record_source.next_record_borrowed()? {
stats.total_records += 1;
progress.log_if_needed(1);

// Push directly to buffer - key extracted inline from raw bytes
buffer.push_coordinate(record.as_ref())?;
buffer.push_coordinate(record)?;

if probe.should_sample_read(stats.total_records) {
probe.log_mid_read(probe_stats(&buffer), Some(pool.phase1_queue_depths()));
Expand Down Expand Up @@ -3020,14 +3024,16 @@ impl RawExternalSorter {
buffer.push(bam_bytes, K::from_full(&full))?;
}

for record in record_source.by_ref() {
// Borrow each record's bytes in place (see the coordinate ingest loop);
// the key is extracted and the bytes copied into the buffer before the
// borrow ends, so no owned `RawRecord` is needed here.
while let Some(bam_bytes) = record_source.next_record_borrowed()? {
stats.total_records += 1;
progress.log_if_needed(1);

// Extract the full template key, verify the lanes the chosen variant
// dropped are constant relative to the first record, then push the
// narrowed key.
let bam_bytes = record.as_ref();
let full = extract_template_key_inline(bam_bytes, lib_lookup, self.cell_tag, cb_hasher);
if let Some(violation) = verify_dropped_lanes(&first, &full, variant) {
let name = String::from_utf8_lossy(
Expand Down
Loading
Loading