Repository navigation
fix(sort): honour --compression-level for the tail of a spilling sort - #653
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 (3)
WalkthroughRisk is incorrect phase ordering during BAM/BAI finalization, which can route output compression through temporary settings. The fix releases merge sources, finalizes output in Phase 2, then deactivates; integration tests compare output and index identity across compression settings. ChangesPhase 2 merge output lifecycle
Compression control and verification
Estimated code review effort: 3 (Moderate) | ~25 minutes Possibly related PRs
Sequence Diagram(s)sequenceDiagram
participant MergeChunks
participant Phase2Guard
participant PooledBamWriter
participant SortWorkerPool
MergeChunks->>Phase2Guard: release_sources()
Phase2Guard->>SortWorkerPool: keep Phase 2 active
MergeChunks->>PooledBamWriter: finish() or finish_index()
MergeChunks->>Phase2Guard: deactivate()
Phase2Guard->>SortWorkerPool: set LEGACY
🚥 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 #653 +/- ##
==========================================
+ Coverage 93.63% 93.68% +0.05%
==========================================
Files 175 175
Lines 107217 107601 +384
==========================================
+ Hits 100388 100803 +415
+ Misses 6829 6798 -31 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@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/external.rs`:
- Around line 3684-3687: The spilling temp-compression integration test does not
cover the indexed finalization path in merge_chunks_with_index. Extend that test
with write_index(true), verify the resulting .bai file exists, and validate that
its virtual offsets resolve correctly after release_sources(), finish_index(),
and deactivate(); otherwise document why the unindexed case is sufficient.
🪄 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: 96615585-b6e3-4d2b-ba78-3deb0fac183c
📒 Files selected for processing (2)
crates/fgumi-sort/src/external.rscrates/fgumi-sort/tests/integration/test_output_compression.rs
ecda4d0 to
eb0d818
Compare
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
`Phase2Guard::deactivate` returned the pool to `LEGACY` *before* the output
writer was finished, in all four merge sites. But `get_sort_priorities` only
offers `CompressOutput` in Phase 2 — in `LEGACY` the sole compress step is
`CompressSpill`, which binds `worker.compressor` (the `temp_compression` one).
So every output block still sitting in the compress queue at teardown, plus the
writer's final flush, was compressed at the temp level and the tail of the BAM
silently ignored `--compression-level`.
The stale comment at both merge sites ("workers still need to compress",
"CompressOutput is eligible whenever the queue is non-empty") describes the
intent correctly; `is_step_eligible` does gate `CompressOutput` on a non-empty
queue. What it misses is that in `LEGACY` the step is never dispatched in the
first place, so the eligibility check is moot.
Measured on a 5,000-pair coordinate sort forced to spill, writing at
`--compression-level 0`:
temp_compression 0 -> 2,062,091 bytes
temp_compression 9 -> 2,021,843 bytes
~40 KB (~2%) of the output came out at the temp level. The in-memory path was
already correct: it never builds a `Phase2Guard`, so it has no early teardown,
and it produces exactly the 2,062,091-byte level-0 output either way. That the
merge path was the odd one out is the argument that this fix, though it does
not touch the underlying coupling, is at a defensible depth.
Give the guard a `finish_output` combinator that owns the whole teardown
order: it releases the drained merge sources, runs the caller's finalizer, and
only then returns the pool to `LEGACY`. Releasing the sources first is what
keeps a slow `finish` from holding the merge's file descriptors and per-file
reorder buffers alive while it drains; finalizing before the phase reset is
what keeps the trailing blocks routed to `output_compressor`. `Drop` still
deactivates, so error paths are unchanged.
The ordering is the whole fix, so it is expressed as one method rather than as
a pair of calls each of the four sites has to sequence correctly — getting it
wrong produces a wrongly-compressed BAM tail, not an error. The two
`initial_keys.is_empty()` branches build their writer inside the closure,
which also closes the window where the writer was constructed outside Phase 2.
The root cause is one level down: which compressor a block gets is decided by
the pool's *global phase* when a worker pops it, rather than by the block's own
origin, even though the single submit site knows which it is. Tagging each
`CompressJob` with its target would retire this ordering contract and the
mirror-image one on `set_phase` (Phase 1 spill jobs still queued when the phase
flips are picked up by `CompressOutput`); that is a separate change, noted in
`finish_output`'s docs.
The new integration test sorts twice at `output_compression 0`, varying only
the temp level, and asserts the two outputs are the same size. It fails on all
three pooled sort orders before this change.
The integration test also gains a `coordinate_indexed` case, so the spilling
merge runs through `merge_chunks_with_index` as well, where the tail block's
level additionally moves the BGZF boundaries the BAI records. That case asserts
the sidecar exists, that an indexed whole-reference query resolves to exactly
the record stream a linear scan of the same file yields, and that the two temp
levels produce byte-identical indexes.
The `initial_keys.is_empty()` early return in each merge finalizes through the
same `finish_output` and so has the same way to go wrong, but Phase 1 checks
its spill trigger only after pushing a record and therefore never writes an
empty chunk — the branch is unreachable from `sort()`. Four unit tests drive
it directly from an empty BGZF chunk: one per merge function, plus two on
`Phase2Guard` itself. The first asserts the phase from *inside* the finalizer,
which is the only place the invariant is observable — asserting either side of
it passes even when the output is finalized in `LEGACY`. The second covers the
error path, pinning the reason `finish_output` binds the finalizer's result
rather than applying `?` to it: the phase reset must not depend on `Drop`.
Reading the phase back needs `SortWorkerPool::current_phase`, a test-only
accessor, since `shared` is private to `worker_pool`.
eb0d818 to
e0c4f26
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.
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.
…ase (#658) 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.
Follow-up to #633. That PR fixed the fully-in-memory sort, which never entered Phase 2 at all and so wrote its entire output at
temp_compression. This fixes the remaining leak on the spilling path, where only the tail of the file is affected.The bug
Phase2Guard::deactivatereturned the pool toLEGACYbefore the output writer was finished, at all four merge sites. Butselect_prioritiesonly offersCompressOutputin Phase 2 — inLEGACYthe sole compress step isCompressSpill, which bindsworker.compressor(thetemp_compressionone). Every output block still sitting in the compress queue at teardown, plus the writer's final flush, was therefore compressed at the temp level, and the tail of the BAM silently ignored--compression-level.The comments at both merge sites ("workers still need to compress", "CompressOutput is eligible whenever the queue is non-empty") describe the intent correctly, and
is_availablereally does gateCompressOutputon a non-empty queue. What they miss is that inLEGACYthe step is never dispatched in the first place, so the eligibility check never gets consulted.Measured
5,000-pair coordinate sort forced to spill, writing at
--compression-level 0:temp_compression~40 KB (~2%) of the output came out at the temp level. The in-memory path is unaffected — it never builds a
Phase2Guard, so there is no early teardown to get wrong, and it produces exactly the 2,062,091-byte level-0 output either way.The fix
Split the guard's teardown into its two halves:
release_sourcesdrops the consumer and unpublishes the spill files — which is what stops workers touching files that are about to be deleted — while leaving the pool in Phase 2, so trailing blocks still route tooutput_compressor.deactivatereturns the pool toLEGACYafter the writer is finished, and remains whatDropcalls, so error paths are unchanged.try_phase2_file_workreturns immediately on an empty file vector, so output compression is the only Phase 2 work left in that window.Testing
New integration test
test_temp_compression_does_not_reach_the_output_bamsorts twice atoutput_compression 0, varying only the temp level, and asserts the two outputs are the same size. It fails on all three pooled sort orders before this change:sort_at_levelgained atemp_compressionparameter so both tests share it; the existing test passes1as before.cargo ci-fmt,cargo ci-lintclean;cargo ci-test6476 passed / 27 skipped.Summary by CodeRabbit