refactor(sort): pick the compressor from the job, not the pipeline phase - #658
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (6)
WalkthroughBGZF compression selection now comes from each queued job’s ChangesCompression routing
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant PooledWriter
participant StagingBuffer
participant SortWorkerPool
participant BGZFCompressor
PooledWriter->>StagingBuffer: assign Spill or Output target
StagingBuffer->>SortWorkerPool: enqueue CompressJob
SortWorkerPool->>BGZFCompressor: compress using job.target
Possibly related PRs
Suggested reviewers: 🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #658 +/- ##
==========================================
+ Coverage 93.80% 93.91% +0.11%
==========================================
Files 177 177
Lines 107658 107779 +121
==========================================
+ Hits 100990 101224 +234
+ Misses 6668 6555 -113 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/fgumi-sort/src/worker_pool.rs`:
- Around line 2200-2214: Add a test near the worker-pool compression tests that
constructs a SortWorkerPool with Zstd, submits a CompressJob using
SpillCodec::Zstd and CompressTarget::Output through submit_compress, and asserts
the worker panics with the existing “zstd is spill-only” message. Ensure the
test reliably allows the worker to process the job while preserving the
unconditional invariant in the Zstd arm.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 57e2bbdb-724e-436a-a51a-5992937d345e
📒 Files selected for processing (6)
crates/fgumi-sort/src/bgzf_io.rscrates/fgumi-sort/src/external.rscrates/fgumi-sort/src/pooled_bam_writer.rscrates/fgumi-sort/src/pooled_chunk_writer.rscrates/fgumi-sort/src/worker_pool.rscrates/fgumi-sort/tests/integration/test_output_compression.rs
05c0001 to
1d38610
Compare
|
@coderabbitai review |
✅ Action performedReview finished.
|
One `compress_queue` carries both spill and output blocks, and `CompressJob`
recorded only the *codec* (BGZF vs zstd), never which of the two it was. So a
worker popping a job re-derived the compression level from `shared.phase`, a
mutable atomic that advances independently of what is already sitting in the
queue. Every block that outlived a phase transition was compressed at the wrong
level, silently.
That single coupling produced the same bug twice, in opposite directions:
- Output blocks popped outside Phase 2 went through the spill compressor at
`temp_compression`, so `--compression-level` was silently discarded. Fixed
twice at the symptom layer — by entering Phase 2 even for an in-memory sort
(#633), then by keeping the pool in Phase 2 across the writer's `finish()`
(#653).
- Phase 1 spill blocks still queued when `begin_phase2` fired went through the
output compressor at `output_compression`. This one was never a live bug: it
was held off only by `drain_pending_spill` happening to sit immediately
before `enter_output_phase`, and was written down as a caller obligation on
`begin_phase2` rather than enforced. Its cost is wasted CPU and disk, not
wrong bytes, since any BGZF level decodes.
The identity was already known for certain at submit time, by the one caller who
cannot be wrong about it: there is a single production submit site,
`StagingBuffer::flush`, and its two constructors are statically distinct — the
output writer and the spill writer. So record it rather than infer it.
`CompressJob` gains a `target: CompressTarget { Spill, Output }` beside the
existing `codec`; `StagingBuffer` carries it and stamps every job it submits;
`try_compress` selects `compressor` or `output_compressor` from `job.target`.
Nothing about the choice reads shared mutable state, so a block is compressed at
the level its writer asked for however long it waited in the queue and whatever
phase the pool is in when a worker gets to it.
This is the second form of that selection. `5df7e176`, which introduced the
worker pool, read it from `shared.phase` at pop time and then narrowed it to the
dispatched `SortStep` to close the window where `set_phase` fired between the pop
and the choice. But a step is itself chosen from the phase, so blocks that
outlived a transition were still mis-compressed — the window was narrowed, not
closed.
With the level travelling on the job, `SortStep::CompressSpill` and
`CompressOutput` did the same thing, so they collapse into one `Compress`. There
was only ever one queue; two steps implied two, and the per-step stats buckets
they fed ("CmpSpl" / "CmpOut") could not be trusted to mean what they said.
`SortStep::COUNT` drops from 5 to 4. `SortStep` is `pub` but `worker_pool` is
`pub(crate)`, so this is not an API change.
Also drops `begin_phase2`'s caller obligation, which no longer exists, and pins
zstd's spill-only invariant now that it is expressible: there is one zstd
compressor per worker, fixed at `temp_compression`, because the output BAM is
always BGZF. That one asserts unconditionally rather than in debug only — an
`Output` job reaching it would be the same silent wrong-level output this
mechanism exists to prevent, and the check is one compare inside the zstd arm,
which the BGZF path never evaluates.
The ordering #653 introduced — release the drained sources, finalize, then
leave Phase 2 — is kept, but for the reason that survives: not holding the
merge's file descriptors and 2 MiB-per-chunk reorder buffers across a slow
`finish`. Its compression rationale is gone, and its docs and tests say so.
Reverting that ordering now leaves `test_temp_compression_does_not_reach_the_
output_bam` passing, which is the check that this commit fixes the cause rather
than a third symptom.
`test_compress_target_decides_level_regardless_of_phase` sweeps the phase across
`LEGACY`/`PHASE1`/`PHASE2` and asserts a `Spill` job stays stored at
`temp_compression = 0` while an `Output` job compresses at `output_compression =
9`. All three cases fail if the compressor is selected from the phase, including
the mirror direction that had no coverage before.
No measurable throughput cost. The job grows by one byte in a struct that
already owns a `Vec<u8>`, and the per-block work is one `match` on a `Copy` enum
replacing one `match` on the step; scheduling is untouched, since the two
collapsed steps shared an eligibility predicate and neither was ever exclusive
(only `ReadInputBlocks` is).
Measured on a spilling coordinate sort of a 3,689,310-record BAM (idt-cfdna,
162 MB, 4 spill chunks, 8 threads, `-m 256m`, output level 6, temp level 1),
10 interleaved reps per binary, baseline being #653's merge commit (`ef0f6f8b`,
identical content in these files):
CPU (user+sys) baseline 11.497s ± 0.312 candidate 11.326s ± 0.462 -1.5%
wall baseline 5.663s ± 0.941 candidate 5.910s ± 1.044 +4.4%
Welch t on CPU time is -0.97, so the difference is not significant; wall clock
on this host is too noisy (σ ≈ 1 s) to resolve anything under ~15%, which is why
CPU time is the reported metric. All 20 runs exited 0 and wrote exactly
3,689,310 records. Byte-for-byte, the two binaries produce the same output: the
sorted BAMs differ only in the `@PG` `CL:` field (it records the binary's path),
and the record streams have identical checksums.
1d38610 to
72d0e2b
Compare
|
@coderabbitai review |
✅ Action performedReview finished.
|
Follow-on to #633 and #653, which each fixed one symptom of the same root cause. This fixes the cause.
The coupling
One
compress_queuecarries both spill and output blocks, andCompressJobrecorded only the codec (BGZF vs zstd) — never which of the two kinds it was. So a worker popping a job re-derived the compression level fromshared.phase, a mutable atomic that advances independently of what is already sitting in the queue. Every block that outlived a phase transition was compressed at the wrong level, silently.That one coupling produced the same bug twice, in opposite directions:
temp_compression, so--compression-levelwas silently discarded. Fixed twice at the symptom layer: by entering Phase 2 even for an in-memory sort (fix(sort): honour --compression-level when the sort fits in memory #633), then by keeping the pool in Phase 2 across the writer'sfinish()(fix(sort): honour --compression-level for the tail of a spilling sort #653).begin_phase2fired went through the output compressor atoutput_compression. This was never a live bug — it was held off only bydrain_pending_spillhappening to sit immediately beforeenter_output_phase, and was written down as a caller obligation onbegin_phase2rather than enforced. Its cost is wasted CPU and disk, not wrong bytes, since any BGZF level decodes.The fix
The identity was already known for certain at submit time, by the one caller who cannot be wrong about it: there is a single production submit site,
StagingBuffer::flush, and its two constructors are statically distinct — the output writer and the spill writer. So record it rather than infer it.CompressJobgains atarget: CompressTarget { Spill, Output }beside the existingcodec;StagingBuffercarries it and stamps every job it submits;try_compressselectscompressororoutput_compressorfromjob.target. Nothing about the choice reads shared mutable state, so a block is compressed at the level its writer asked for however long it waited in the queue and whatever phase the pool is in when a worker gets to it.This is the second form of that selection.
5df7e176, which introduced the worker pool, read it fromshared.phaseat pop time and then narrowed it to the dispatchedSortStepto close the window whereset_phasefired between the pop and the choice. But a step is itself chosen from the phase, so blocks that outlived a transition were still mis-compressed — the window was narrowed, not closed.With the level travelling on the job,
SortStep::CompressSpillandCompressOutputdid the same thing, so they collapse into oneCompress. There was only ever one queue; two steps implied two, and the per-step stats buckets they fed (CmpSpl/CmpOut) could not be trusted to mean what they said.SortStep::COUNTdrops from 5 to 4.SortStepispubbutworker_poolispub(crate), so this is not an API change.Also drops
begin_phase2's caller obligation, which no longer exists, and pins zstd's spill-only invariant now that it is expressible: there is one zstd compressor per worker, fixed attemp_compression, because the output BAM is always BGZF. That one asserts unconditionally rather than in debug only — anOutputjob reaching it would be the same silent wrong-level output this mechanism exists to prevent, and the check is one compare inside the zstd arm, which the BGZF path never evaluates.The ordering #653 introduced — release the drained sources, finalize, then leave Phase 2 — is kept, but for the reason that survives: not holding the merge's file descriptors and 2 MiB-per-chunk reorder buffers across a slow
finish. Its compression rationale is gone, and its docs and tests say so.Testing
Reverting #653's ordering now leaves
test_temp_compression_does_not_reach_the_output_bampassing — that is the check that this fixes the cause rather than adding a third symptom fix.test_compress_target_decides_level_regardless_of_phasesweeps the phase acrossLEGACY/PHASE1/PHASE2and asserts aSpilljob stays stored attemp_compression = 0while anOutputjob compresses atoutput_compression = 9. All three cases fail if the compressor is selected from the phase, including the mirror direction that had no coverage before.cargo ci-fmt,cargo ci-lint,cargo ci-doc(with-D warnings, since this adds intra-doc links), andcargo ci-test(6597 tests) all pass.Performance
No measurable cost. The job grows by one byte in a struct that already owns a
Vec<u8>, and the per-block work is onematchon aCopyenum replacing onematchon the step. Scheduling is untouched: the two collapsed steps shared an eligibility predicate, and neither was ever exclusive (onlyReadInputBlocksis).Measured on a spilling coordinate sort of a 3,689,310-record BAM (idt-cfdna, 162 MB, 4 spill chunks, 8 threads,
-m 256m, output level 6, temp level 1), 10 interleaved reps per binary, baseline being #653's merge commitef0f6f8b:Welch t on CPU time is −0.97, so the difference is not significant. Wall clock on this host is too noisy (σ ≈ 1 s) to resolve anything under ~15%, which is why CPU time is the reported metric. All 20 runs exited 0 and wrote exactly 3,689,310 records. The two binaries produce equivalent output: the sorted BAMs differ only in the
@PGCL:field, which records the binary's path, and the record streams have identical checksums.Reading order
crates/fgumi-sort/src/worker_pool.rsis the change —CompressTarget, theCompressJobfield,SortStep, andtry_compress. Everything else follows from it:bgzf_io.rsthreads the field, the two pooled writers each stamp their constant, andexternal.rsplus the integration test are documentation that had asserted the old coupling.Summary by CodeRabbit
Bug Fixes
Tests
Documentation