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
52 changes: 46 additions & 6 deletions crates/fgumi-sort/src/bgzf_io.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,9 @@
//! accumulating data into ~64KB blocks before submitting compression jobs.

use crate::codec::SpillCodec;
use crate::worker_pool::{BufferPool, CompressJob, CompressResult, PermitPool, SortWorkerPool};
use crate::worker_pool::{
BufferPool, CompressJob, CompressResult, CompressTarget, PermitPool, SortWorkerPool,
};
use anyhow::Result;
use crossbeam_channel::{Receiver, Sender};
use fgumi_bgzf::{BGZF_EOF, BGZF_MAX_BLOCK_SIZE};
Expand Down Expand Up @@ -39,6 +41,10 @@ pub(crate) struct StagingBuffer {
result_tx: Sender<CompressResult>,
permit_pool: Arc<PermitPool>,
codec: SpillCodec,
/// Stamped onto every job this buffer submits, so the worker that pops it
/// compresses at the level this writer asked for regardless of the pool's
/// phase at that moment.
target: CompressTarget,
}

impl StagingBuffer {
Expand All @@ -49,6 +55,7 @@ impl StagingBuffer {
result_tx: Sender<CompressResult>,
permit_pool: Arc<PermitPool>,
codec: SpillCodec,
target: CompressTarget,
) -> Self {
Self {
pool,
Expand All @@ -57,6 +64,7 @@ impl StagingBuffer {
result_tx,
permit_pool,
codec,
target,
}
}

Expand Down Expand Up @@ -123,6 +131,7 @@ impl StagingBuffer {
serial,
result_tx: self.result_tx.clone(),
codec: self.codec,
target: self.target,
});
Ok(())
}
Expand Down Expand Up @@ -331,7 +340,13 @@ mod tests {
io_writer_loop(writer, result_rx, buffer_pool, pp, codec, None)
});

let mut staging = StagingBuffer::new(Arc::clone(&pool), result_tx, permit_pool, codec);
let mut staging = StagingBuffer::new(
Arc::clone(&pool),
result_tx,
permit_pool,
codec,
CompressTarget::Spill,
);
staging.write_chunked(data).unwrap();
staging.flush().unwrap();
drop(staging); // closes result_tx senders → io_writer_loop exits
Expand All @@ -352,7 +367,13 @@ mod tests {
let (result_tx, _result_rx) = pool.compress_result_channel();
let permit_pool = make_permit_pool(&pool);

let mut staging = StagingBuffer::new(Arc::clone(&pool), result_tx, permit_pool, codec);
let mut staging = StagingBuffer::new(
Arc::clone(&pool),
result_tx,
permit_pool,
codec,
CompressTarget::Spill,
);
// Flush with empty buffer: should not submit a compress job
staging.flush().unwrap();

Expand All @@ -373,7 +394,13 @@ mod tests {
let pool = Arc::new(SortWorkerPool::new(1, 1, 6, codec));
let (result_tx, _result_rx) = pool.compress_result_channel();
let permit_pool = make_permit_pool(&pool);
let mut staging = StagingBuffer::new(Arc::clone(&pool), result_tx, permit_pool, codec);
let mut staging = StagingBuffer::new(
Arc::clone(&pool),
result_tx,
permit_pool,
codec,
CompressTarget::Spill,
);

assert!(!staging.is_full(), "empty buffer should not be full");
staging.buf().extend(vec![0u8; BGZF_MAX_BLOCK_SIZE]);
Expand Down Expand Up @@ -404,7 +431,13 @@ mod tests {
io_writer_loop(writer, result_rx, buffer_pool, pp, codec, None)
});

let mut staging = StagingBuffer::new(Arc::clone(&pool), result_tx, permit_pool, codec);
let mut staging = StagingBuffer::new(
Arc::clone(&pool),
result_tx,
permit_pool,
codec,
CompressTarget::Spill,
);
staging.write_chunked(&large).unwrap();
staging.flush().unwrap();
drop(staging);
Expand Down Expand Up @@ -452,9 +485,16 @@ mod tests {
serial: 1,
result_tx: result_tx.clone(),
codec,
target: CompressTarget::Spill,
});
permit_pool.acquire().unwrap();
pool.submit_compress(CompressJob { data: data1, serial: 0, result_tx, codec });
pool.submit_compress(CompressJob {
data: data1,
serial: 0,
result_tx,
codec,
target: CompressTarget::Spill,
});

// Wait for both compress results to be received by io_writer_loop
io_handle.join().unwrap().unwrap();
Expand Down
107 changes: 55 additions & 52 deletions crates/fgumi-sort/src/external.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1679,32 +1679,25 @@ impl<K: RawSortKey + 'static> Phase2Guard<'_, K> {

/// Release the drained merge sources, finalize the output, then leave Phase 2.
///
/// **The output must be finalized inside Phase 2.** Phase 2 is what puts
/// [`SortStep::CompressOutput`](crate::worker_pool::SortStep::CompressOutput)
/// on a worker's priority list, and that step is what selects
/// `output_compressor`. In `LEGACY` the only compress step workers schedule
/// is `CompressSpill`, so every output block still queued — plus the
/// writer's final flush — would be compressed at `temp_compression`, and the
/// tail of the BAM would silently ignore the caller's `output_compression`.
/// This method exists so that ordering cannot be got wrong at a call site:
/// `finish` runs between the source release and the return to `LEGACY`.
/// Releasing the sources first is the point: it keeps a slow `finish` from
/// holding the merge's file descriptors and per-file reorder buffers (a 2 MiB
/// `BufReader` per spill file) alive while it drains. Leaving Phase 2 last
/// then costs nothing, and keeps the pool's phase an accurate description of
/// what the pool is doing for as long as it is doing it.
///
/// Releasing the sources first is what keeps a slow `finish` from holding
/// the merge's file descriptors and per-file reorder buffers (a 2 MiB
/// `BufReader` per spill file) alive while it drains.
///
/// The underlying coupling — "which compressor a block gets" being a
/// function of the pool's *global phase* rather than of the block's own
/// origin — is what makes this ordering load-bearing at all. Tagging each
/// `CompressJob` with its target at submit time would retire this contract
/// (and the mirror-image one on `set_phase`); that is a separate change.
/// Note what this ordering is *not* responsible for. The output's compression
/// level does not depend on it: each block carries its own
/// [`CompressTarget`](crate::worker_pool::CompressTarget), stamped by the
/// writer that produced it, so trailing output blocks are compressed at
/// `output_compression` whatever phase the pool is in when a worker pops
/// them. That was not always true — the compressor used to be selected from
/// the pool's phase at pop time, which is why finalizing after the return to
/// `LEGACY` silently wrote the tail of the BAM at `temp_compression`.
///
/// Takes `self` by value: finalizing the output is the guard's terminal
/// operation, so a second call — which would run `finish` with the pool
/// already back in `LEGACY`, silently compressing the output at
/// `temp_compression` and reporting success — is a compile error rather than
/// a runtime check. The remaining `Drop` at the end of this method is a
/// no-op, since `deactivate` has already cleared `active`.
/// operation, so a second call is a compile error rather than a runtime
/// check. The `Drop` at the end of this method is a no-op, since
/// `deactivate` has already cleared `active`.
fn finish_output<T>(mut self, finish: impl FnOnce() -> Result<T>) -> Result<T> {
self.release_sources();
let finished = finish();
Expand Down Expand Up @@ -1837,14 +1830,18 @@ impl RawExternalSorter {
///
/// Every sort path calls this exactly once, immediately after ingest
/// completes, to hand the pool over to Phase 2 and lift the active-worker cap
/// to `merge_threads`. Centralized here so the invariant cannot drift between
/// sort orders: it is precisely the transition whose omission left the
/// in-memory output writing at the wrong compression level.
/// to `merge_threads`. Centralized here so the transition cannot drift
/// between sort orders.
///
/// Workers idled by the Phase-1 cap re-check within `MAX_BACKOFF_US`, so no
/// explicit wake is needed; capped workers always drained their held items, so
/// raising the cap can never be required for correctness, only for width. The
/// phase half is what routes output blocks to `output_compressor` — see
/// raising the cap can never be required for correctness, only for width.
///
/// Omitting this call once left in-memory sorts writing their output at
/// `temp_compression`, because the phase used to decide which compressor a
/// block went through. It no longer does — blocks carry their own
/// [`CompressTarget`](crate::worker_pool::CompressTarget) — so what the
/// phase buys now is Phase 2 scheduling and the wider worker cap. See
/// [`SortWorkerPool::begin_phase2`].
fn enter_output_phase(&self, pool: &SortWorkerPool) {
pool.begin_phase2(self.phase2_threads());
Expand Down Expand Up @@ -3491,16 +3488,17 @@ impl RawExternalSorter {
/// single-sourced and the two paths cannot drift. Returns the sources plus
/// an RAII [`Phase2Guard`] (borrowing `pool`); callers finish through
/// [`Phase2Guard::finish_output`], and `Drop` resets Phase 2 on any path
/// that does not get that far. Phase 2 is armed whether or not there are disk
/// chunks: it is what makes workers reach for `output_compressor`
/// (`SortStep::CompressOutput`) rather than the Phase 1 spill compressor, so
/// leaving an all-in-memory merge in Phase 1 would silently write the output
/// BAM at `temp_compression` and discard the caller's `output_compression`.
/// With no spill files `try_phase2_file_work` sees an empty file vector and
/// returns immediately, so the only Phase 2 work available is output
/// compression. The consumer is created only for disk sources — it drains the
/// pool's per-file reorder buffers, and there is nothing to drain when
/// everything is already in memory.
/// that does not get that far.
///
/// Phase 2 is armed whether or not there are disk chunks: it is what lets
/// workers schedule `SortStep::Phase2FileWork` and reach the merge's
/// per-file reorder buffers. With no spill files `try_phase2_file_work` sees
/// an empty file vector and returns immediately, so the only work left is
/// compression — which is correct in any phase, since each block carries its
/// own [`CompressTarget`](crate::worker_pool::CompressTarget). The consumer
/// is created only for disk sources — it drains the pool's per-file reorder
/// buffers, and there is nothing to drain when everything is already in
/// memory.
fn setup_phase2_merge<'a, K: RawSortKey + Default + 'static>(
chunk_files: &[PathBuf],
memory_chunks: MemorySources<K>,
Expand All @@ -3518,8 +3516,10 @@ impl RawExternalSorter {
let sources = Self::build_chunk_sources::<K>(chunk_files, memory_chunks, pool)?;
debug!("Merging from {} sources...", sources.len());

// Activate Phase 2 for the whole merge, spill files or not, so the output
// BAM is compressed at `output_compression` rather than `temp_compression`.
// Activate Phase 2 for the whole merge, spill files or not, so workers can
// schedule `SortStep::Phase2FileWork` and reach the per-file reorder
// buffers. Compression level is no longer a reason to arm it — each block
// carries its own `CompressTarget` and is compressed correctly in any phase.
// The consumer is created only for disk sources: it holds an Arc snapshot of
// the pool's per-file Phase 2 state — workers populate per-file reorder
// buffers, the consumer pops from them. There is nothing to consume when
Expand Down Expand Up @@ -4827,19 +4827,22 @@ mod tests {
}
}

/// The output must be finalized while the pool is still in Phase 2.
/// The merge sources are released *before* the output is finalized, and the
/// pool leaves Phase 2 only *after*.
///
/// This is the invariant [`Phase2Guard::finish_output`] exists to enforce,
/// and the one the bug violated: the trailing output blocks are compressed
/// while the writer is being finished, and `CompressOutput` — the step that
/// selects `output_compressor` — is only on a worker's priority list in
/// Phase 2. Finalizing in `LEGACY` sends those blocks through the spill
/// compressor at `temp_compression`.
/// [`Phase2Guard::finish_output`] exists for the first half: a `finish` that
/// blocks draining the writer must not still be pinning the merge's file
/// descriptors and 2 MiB-per-chunk reorder buffers. So the observation is
/// made from *inside* the finalize closure, where the real writer's
/// `finish()` runs — not merely before and after it.
///
/// So the assertion is made from *inside* the finalize closure, where the
/// real writer's `finish()` would run — not merely before and after it.
/// The phase is checked in the same place, but note it is no longer
/// load-bearing for compression: blocks carry their own `CompressTarget`, so
/// the tail is written at `output_compression` in any phase. What it pins is
/// that the pool's phase still describes what the pool is doing while it is
/// doing it.
#[test]
fn test_phase2_guard_finishes_output_while_still_in_phase2() {
fn test_phase2_guard_releases_sources_before_finalizing_output() {
use crate::keys::RawCoordinateKey;
use crate::worker_pool::phase;

Expand All @@ -4862,8 +4865,8 @@ mod tests {
assert_eq!(
observed,
(phase::PHASE2, 0),
"the output must be finalized in Phase 2 (so its blocks are compressed at \
output_compression) with the drained spill files already unpublished"
"the output must be finalized with the drained spill files already unpublished, and \
before the pool leaves Phase 2"
);
// `finish_output` consumed the guard, so its `Drop` has already run — a
// second deactivate must not have disturbed the phase it just set.
Expand Down
12 changes: 10 additions & 2 deletions crates/fgumi-sort/src/pooled_bam_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@

use crate::bgzf_io::{BlockOffset, StagingBuffer, io_writer_loop};
use crate::codec::SpillCodec;
use crate::worker_pool::{CompressResult, PermitPool, SortWorkerPool};
use crate::worker_pool::{CompressResult, CompressTarget, PermitPool, SortWorkerPool};
use anyhow::Result;
use crossbeam_channel::{Receiver, bounded, unbounded};
use fgumi_bam_io::BaiBuilder;
Expand Down Expand Up @@ -124,7 +124,15 @@ impl PooledBamWriter {
io_writer_loop(writer, result_rx, buffer_pool, pp, SpillCodec::Bgzf, block_offset_tx)
});

let mut staging = StagingBuffer::new(pool, result_tx, permit_pool, SpillCodec::Bgzf);
// `CompressTarget::Output`: every block this writer submits is the sort's
// output BAM and must be compressed at `output_compression`.
let mut staging = StagingBuffer::new(
pool,
result_tx,
permit_pool,
SpillCodec::Bgzf,
CompressTarget::Output,
);

// Write BAM header into a temporary buffer then flush in BGZF-sized chunks.
// Headers can exceed BGZF_MAX_BLOCK_SIZE; write_chunked handles the splitting.
Expand Down
12 changes: 10 additions & 2 deletions crates/fgumi-sort/src/pooled_chunk_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
use crate::bgzf_io::{StagingBuffer, io_writer_loop};
use crate::codec::{SpillCodec, ZSPILL_MAGIC};
use crate::keys::RawSortKey;
use crate::worker_pool::{CompressResult, PermitPool, SortWorkerPool};
use crate::worker_pool::{CompressResult, CompressTarget, PermitPool, SortWorkerPool};
use anyhow::Result;
use crossbeam_channel::bounded;
use fgumi_bgzf::BGZF_MAX_BLOCK_SIZE;
Expand Down Expand Up @@ -78,7 +78,15 @@ impl<K: RawSortKey> PooledChunkWriter<K> {
thread::spawn(move || io_writer_loop(writer, result_rx, buffer_pool, pp, codec, None));

Ok(Self {
staging: Some(StagingBuffer::new(pool, result_tx, permit_pool, codec)),
// `CompressTarget::Spill`: every block this writer submits is a
// Phase 1 spill chunk and must be compressed at `temp_compression`.
staging: Some(StagingBuffer::new(
pool,
result_tx,
permit_pool,
codec,
CompressTarget::Spill,
)),
key_buf: Vec::new(),
io_handle: Some(io_handle),
_phantom: PhantomData,
Expand Down
Loading
Loading