Repository navigation
feat(pipeline): typed-step pipeline foundation - #870
Conversation
…672) First of a crate-by-crate series landing the runall / unified-pipeline work. Adds the two dependency-free leaf crates, wired into the workspace but not yet called by any command. fgumi-fmt hosts the pure formatting helpers -- format_count, format_duration and format_rate -- with no dependencies at all, so any layer can reach them without inheriting clap or sysinfo. fgumi-cli-common holds the shared CLI types and re-exports those three; it depends only on fgumi-fmt. Both crates arrived with per-crate version pins; those are converted to [workspace.dependencies] entries to match the workspace's existing centralization. That required promoting enum_dispatch, num_cpus and sysinfo, each now shared by two members. sysinfo matters most: it is held at 0.38 because >= 0.39 requires rustc 1.95, above the pinned 1.93 MSRV, so a per-crate pin could drift and silently raise the effective MSRV. Extends fgumi-cli-common's tests to cover the paths the crate arrived without: the OperationTimer, validate_file_exists, the parse_memory_size error arms, parse_memory/parse_memory_reserve, the binary-unit Display impls, and the non-per-thread and over-commit branches of the memory-budget resolver. Also fixes three defects in the umbrella crate that the incoming copies had already corrected, so the live code and the future replacement agree: - system.rs: detect_total_memory saturated at usize::MAX, which made the downstream `budget > total` guard in resolve_memory_budget unreachable on 32-bit targets. Now saturates at usize::MAX / 2. - validation.rs: the too-large plain-number error suggested `mb_value / 1000` labelled GB. The input is MiB, so the figure was wrong -- 2,000,000 MiB was reported as "2000GB" when the true value is 2097 GB, or 1953 GiB. Now divides by 1024 and labels it GiB. - commands/common.rs: CompressionOptions derived Default, yielding compression_level 0 (uncompressed BGZF) while the clap default is 1, so programmatic callers silently wrote uncompressed BAM. Now hand-written to match. Reachable today only from tests and doc examples. Each of the three is pinned by a test that fails if the fix is reverted.
…te macro (#674) Adds the proc-macro crate runall uses to re-expose every per-stage option of `fgumi sort` / `group` / `simplex` / `duplex` / `codec` without hand-maintaining a parallel option set. `#[multi_options("sort", "Sort Options")]` takes a `clap::Args` options struct and generates a `MultiSortOptions` companion whose flags are named `--sort::<flag>`, plus `validate()`, `TryFrom` and `From` conversions. The companion is faithful to the standalone command by construction: * every `default_value*` attribute is copied verbatim, so the two flags cannot advertise different defaults and the field type needs no `Display` impl; * long aliases are re-prefixed and short flags dropped, so nothing un-namespaced reaches the parent command; * each generated arg carries its own `help_heading`, so a parent argument declared after a flattened companion keeps its own heading; * `#[doc]` attributes are forwarded individually, preserving clap's short/long help split; * `#[cfg]` gates are forwarded onto the field and both conversion arms; * `#[arg(skip)]` fields ride along as skip fields, making `From` + `validate()` a lossless round trip; * `Vec<T>` passes through untouched -- clap never consults a struct's `Default`, so backfilling from it would diverge from the standalone command. Required-ness stays staged: `required` is dropped from the clap side and enforced by `validate()`, so a missing value is reported as `--sort::max-memory is required when sort is selected` rather than by clap's parser, which cannot name the stage. Forms the macro cannot re-expose faithfully fail the build with a spanned `syn::Error` pointing at the offending field: cross-field reference keys whose arg ids would dangle after prefixing, `id`/`name` overrides, clap's `key(value)` call form for classified keys, field- and struct-level `#[command(...)]`, `#[group(...)]`, an unenforceable `required`, a non-identifier prefix, and the legacy `#[clap(...)]` / `#[structopt(...)]` spellings that every classifier would otherwise silently ignore. Tests include parity fixtures carrying the `#[arg(...)]` attributes of the real `SortOptions` and `GroupOptions` verbatim, asserting the standalone and prefixed commands parse to identical structs, plus ten trybuild cases pinning the diagnostics. Also adds fgumi-cli-macros, fgumi-fmt and fgumi-cli-common to publish.yml's CRATES list. The latter two are publishable but were never listed, which fails that workflow's own completeness check on push to main.
… framework (#697) Adds the `Step`/`Step2` traits, the bounded queue layer, the reorder stage and the work-stealing runtime that schedules and drives them, as a leaf crate. Additive: nothing on the branch changes behaviour and the crate has no in-tree consumer yet. The crate holds nothing that reads or writes sequencing data, so it depends only on ahash, crossbeam-queue, parking_lot, log and one noodles::sam type. Changes made while porting: - Dependencies moved onto `[workspace.dependencies]`. `trybuild` gains a second consumer here, so it is centralized too and `fgumi-cli-macros` switches to `{ workspace = true }`. - Five `collapsible_if` errors fixed as let-chains rather than suppressed. - `fgumi-pipeline-core` added to `publish.yml`'s `CRATES` array; it is a leaf, so any position before its future dependents is valid. - The crate-root doc described the crate as an extraction re-exported by `fgumi` as `pipeline::core`, which is true neither now nor after the series lands; rewritten to describe the crate itself. Two references to design docs that exist in no branch are dropped, and three pre-extraction `core/*.rs` module paths become intra-doc links. - Twenty rustdoc link defects fixed: nine unresolved links now carry crate paths, eight public-doc links to private items are demoted to code spans, three redundant explicit link targets shortened. `ci-doc` runs with `-D warnings`; these never surfaced while this was a private module. Review tooling follows the code out of `src/lib/`: - `.coderabbit.yaml`'s pipeline `path_instructions` glob named only the `src/lib/{unified_pipeline,pipeline}` homes, so extracting the engine into a crate silently dropped it from review — including the clause about new `unsafe` in the typed-step dispatch hot path, written for exactly the sites this commit ships. Widened to cover the crate: 33 of the 42 changed files are instructed now, against 4 before. - `miri.yml` gains a step for `fgumi-pipeline-core::erased`, where all four `#[allow(unsafe_code)]` sites live. It runs clean in ~6s. Both Miri steps now list the filtered tests first and fail on an empty selection: `cargo test <filter>` exits 0 when a filter matches nothing, so renaming either module would have retired the gate silently. The rest of the crate is excluded deliberately: its runtime tests spawn threads and take Miri well past ten minutes, and two wedge tests use wall-clock watchdogs that abort under Miri's slowdown. The four `#[allow(unsafe_code)]` sites in `erased.rs` cache the per-dispatch `downcast_ref` of a step's typed handles by widening `&'a` to `&'static` for storage and narrowing it back on read. The invariant the compiler cannot check -- that every dispatch passes the same handle box for a given step -- is enforced on the cached-hit path by storing the erased box's address in the cache slot and asserting it unconditionally -- release included, since a `debug_assert!` alone would leave release builds reading through a stale pointer. The check is one load and one compare, not the `downcast_ref` `TypeId` probe the cache exists to elide; a debug-only `debug_assert!` still re-resolves the typed pointer as a second diagnostic. CLAUDE.md's unsafe-code section gains an entry for them; note it documents all four, where feat-runall's own entry describes only the two `TypedStep` sites and omits the `TypedStep2` pair. 57 tests added over the 284 that arrive with the crate, all against public API that had no coverage: `build_queues`/`mark_all_drained` round-trips for every fan-out output shape (asserting branches are open before and closed after, so they cannot pass on an inert `mark_all_drained`), `Chain::into_multi` for the 3-, 4- and ordered-bytes-2 branch cases, the type-erased `append_source`/`append_step`/`append_step2` assembly API, `HeapSize` for String/Vec/Option, `BuildError` rendering, the `PipelineConfig` setters, `Step2`'s trait defaults, and the two input-arity dispatch guards. Fixes from review, all in the pipeline-core runtime: - `driver.rs`: a sticky step's `Progress` also means "held an item", so one sitting on a full output reported it forever. The sticky fast path is now bounded at `STICKY_BURST_LIMIT` re-entries AND round-robin no longer takes the priority restart on the sticky owner's `Progress` -- with only one of the two, the walk still breaks before reaching the consumer that drains that output, and a one-worker chain livelocks. Both halves are pinned by `sticky_holding_source_yields_to_its_draining_consumer`, which wedges (watchdog abort) if either is reverted. - `erased.rs`: `build_fused_output_set` forced `QueueSpec::Unbounded` on every branch. The fused driver runs a producer before its consumer, so a step emitting more per `try_run` than its consumer removes grew the edge every pass, without limit. The profile's count/byte bound is kept; only the ordering is dropped. - `stats.rs`: the SPIN verdict exempted Detached steps by NAME while the edge attribution beside it deliberately used `StepIdx`. `StatsSnapshot:: detached` now carries the index and the exemption matches on it, so duplicate step names stay independent. - `handles.rs`: `retry` on an `Ordered` branch with `ordinal = None` re-pushed, allocating a fresh ordinal and abandoning the original -- a permanent gap that stalls the consumer's `ReorderStage`. It now panics in every build; the recoverable `Direct` sibling is left as-is. - `contexts.rs`: `find_producer` returned the first matching edge, so two producers wired to one single-input consumer left the second branch unpopped. It now rejects that graph, mirroring the per-slot uniqueness check `find_all_producers` already had. - `tests.rs`: `PairSummer`/`OrderedPairSummer` reported `Finished` while a branch item was still buffered, silently dropping it -- the exact loss their own doc warns step authors about. `run_with_deadlock_timeout` reported a panic inside the run as "DEADLOCKED or stalled"; it now re-raises the original panic and blames a deadlock only on a true timeout. - `compile_fail.rs`: documents that the `.stderr` fixtures pin rustc's whole diagnostic, including its trait-impl inventory, and how to tell a re-bless from a real regression. The inventory blocks cannot just be deleted -- trybuild compares for equality, so that fails the test outright. Three follow-ons to those fixes, found by re-reading the call sites they touched rather than by review: - `queues.rs`: the same `debug_assert!`-with-a-data-losing-release-fallback shape the `handles.rs` fix removed was live in all three transports' `try_push` -- a push after `mark_drained` went into a queue the consumer had already closed, so the item was silently dropped and surfaced only as a short output. The check is now the unconditional `assert_not_drained` helper, shared by `CountBoundedQueue`, `ByteBoundedQueue` and `UnboundedQueue` and naming which transport failed. It loads `Relaxed` rather than `Acquire`: `drained` is monotonic (its only write anywhere is `store(true, Release)`), so a `Relaxed` load can miss a detection under a cross-thread race but can never panic a correct program, and this now sits on the per-item push path. - `fused.rs` / `builder.rs`: with fused transports keeping their profile byte bounds, `--queue-memory-total` was being silently ignored on that path in favour of the per-step defaults -- the fused contexts are built inside `run_fused_single_thread`, so `Pipeline::run` could not apply the budget on their behalf. The budget is now threaded in and applied there; `apply_initial_queue_budget` becomes `pub(crate)` for it. - `builder.rs` / `step.rs`: two comments the above invalidated. The fused fast-path comment justified skipping the rebalancer with "has no bounded queues to rebalance", which stopped being true; it now gives the real reason (producer-then-consumer in one pass, so no cross-worker imbalance) and points at where the budget is applied. `step.rs`'s last-worker-barrier note still said the `try_push`-after-drain violation was "debug-asserted". Third review round: - `handles.rs`: the 2-, 3- and 4-branch fan-out builders used the non-byte-aware `build_branch`, which panics on `QueueSpec::ByteBounded`, so ANY fan-out step declaring a byte bound failed to build. An armed deadlock monitor requires `ByteBounded` on every output transport, so a monitored fan-out chain was unbuildable. All nine calls now use `build_branch_byte_aware`, which delegates non-byte specs straight back to `build_branch` -- count/unbounded branches are unchanged. Three rstest cases assert the declared bound is registered on every branch of every shape; all three panic without the swap. - `sampler.rs`: each tick read every ordered edge's reorder stash TWICE -- once for the histogram, once for the timeline row -- taking the `ReorderStage` state mutex (held by every must-accept push and `try_pop_in_order`) on both, and the two reads could disagree about one tick. Depths are now read once per tick via `read_depths` and shared by both consumers. The module doc claimed "one `Relaxed` load per edge per tick" and "never perturbs the worker hot path", which was never true for ordered edges; it now describes the real cost. - `miri.yml`: both steps now also assert that no `#[allow(unsafe_code)]` site has moved out of the Miri-scoped module. The `--list` guard proves the filter selects tests, not that it still covers the crate's `unsafe` -- a site that moved would leave the filter matching and the site unchecked. - `storage.rs`: `affinity_in_range` re-derived the affinity->worker mapping by matching variants; it now routes through `Affinity::target_worker`, the documented single source of truth, so the range check and the dispatch gate cannot disagree. Added the missing `should_panic` case for the Exclusive owner-range assert -- only the Serial arm was covered, so that always-on `assert!` could have been weakened to a `debug_assert!` with CI still green. - `queues.rs`: the slot-cap reject test filled with 0-byte items, exercising the byte-reservation rollback with `size == 0` where a leak or double-subtract is invisible. Added a nonzero-size case pinning `current_bytes` across the reject. - Test-fixture steps returned `StepOutcome::Contention` when an output push was rejected. The driver treats it like `NoProgress`, but it feeds `contention_count`, from which the bottleneck verdict derives its SPIN ratio -- so ordinary backpressure was inventing mutex contention and could trip a bogus SPIN finding. Fixed in all six sites: the two flagged in `tests.rs` (`OrderedSource`, `OrderedPairSummer`) plus the same shape in `detached.rs` and `builder.rs`, which the review did not name. - Docs corrected where they contradicted the code: `scheduler.rs` claimed sticky steps do not take part in the round-robin walk (they do -- only their `Progress` priority restart is suppressed), `step.rs`'s `Affinity` omitted that `Detached` also ignores the hint and is never range-checked, and `TypedStep2:: build_output_set` rebuilt `profile()` twice. Fourth review round: - `handles.rs`: `BranchInputHandle::pop` reported every popped item's `heap_size()` to `record_pop`, but a `CountBoundedQueue` / `UnboundedQueue` records `record_push(0)` by design (no byte accounting on its push path). So count/unbounded edges carried nonzero `popped_bytes` against zero `pushed_bytes`, and `compute_edge_stats` reported a `mibytes_per_s` throughput for an edge that measures no bytes at all. The handle now carries `record_item_bytes`, gated on the same `transport_handle.is_some()` signal `ordered_branch` already used for the push side. - `queues.rs`: `ByteBoundedQueue::set_limit_bytes` accepted 0, which makes `try_push` reject unconditionally (`current_bytes >= 0`) and wedges the producer forever. `new` asserts `limit_bytes > 0` for exactly that reason, so the setter could undo the constructor's invariant. Now floored at 1 -- clamped rather than asserted because this runs on a live pipeline, where a 1-byte limit still makes progress and a panic would kill the run over a recoverable arithmetic slip. - `detached.rs`: `DetachedPlaceholder::detached_group` returned `PerStep` unconditionally, so a step that declared `DetachedGroup::Shared(..)` read back as `PerStep`. `extract_detached_steps` is `pub` and leaves these placeholders in the caller's slice, so an external caller regrouping from it would split one shared group across separate driver threads. The placeholder now carries the real label. - `detached_two_sided_no_deadlock` asserted only that N items arrived, so a step emitting N copies of one item -- or duplicating the held value while dropping a popped one -- passed. It now collects the values and asserts the sorted multiset (sorted, not sequenced: the producer's backpressure branch legitimately reorders on retry). - `outputs.rs`: the block comment claimed one `build_queues`/`mark_all_drained` round-trip per output shape; `Single<T>`, `OrderedBytesSingle<T>` and the plain `(A, B)` tuple had none -- and `OrderedBytesSingle<T>` is the only shape routing through `build_single_queues_ordered_bytes`, so its drain path had no coverage at all. All three added; every shape with a `StepOutputs` impl now has one. - `held.rs`: `double_put_panics` used `catch_unwind`, which accepts any panic, so a `put` failing for an unrelated reason would still pass while the double-put guard itself had gone. Pinned with `#[should_panic(expected = ...)]`, matching every other always-on guard here. Fifth review round. Two substantive bugs, both in error/stall handling: - `signal.rs`: `to_result` mapped a failed run to `Ok(())`. `record_error` wins the state CAS before it runs `payload.set`, so `state == STATE_ERROR` with `outcome() == None` is reachable; the `None` arm rescued only `STATE_CANCELLED`, so the error case fell through to success. The cancel side of this exact publish window already had a guard and a test -- the error side was the sibling nobody closed. Now derived from the terminal state, which cannot regress if a future driver forgets to join a writer. A successful run never reaches the new arm: `is_done()` is `state != STATE_OK`. - `fused.rs`: the driver failed the run on the FIRST pass in which no step reported progress. `NoProgress` is the transient "input momentarily empty but not drained" outcome -- a source waiting on a background reader returns it -- so a healthy run was aborted and its output truncated, and the `debug_assert!` turned the same transient pass into a panic under test. Idle passes now back off and retry against a wall-clock budget, and only a sustained stall reports the error. Budgeted by time rather than by pass count because `thread::sleep` granularity varies by platform. The loop also had no backoff at all, so it span the calling thread at 100%. Two findings were rejected as false positives, verified rather than assumed: - An ordered non-byte-bounded branch was said to get an uncapped reorder stash via `ReorderStage::new`. All seven `ordered_branch` call sites pass `DEFAULT_REORDER_OVERFLOW_BYTES`, so that arm was unreachable from any branch builder. Rather than only reply, the `Option<u64>` parameter is now a plain `u64` and the dead unbounded arm is gone, so the question cannot be asked of this code again. - The `HeapSize` compile-fail fixture was said to risk "expected failure but succeeded" if the bound lived only on the `StepOutputs` impl. It is on the `OrderedBytesSingle<T>` declaration, and the fixture's pinned `.stderr` shows the `E0277` plus the `required by a bound in OrderedBytesSingle` note. The rest were accuracy and test-strength fixes: - The cached-handle address assert was documented as though it proved box identity. It compares data addresses only, so allocator address reuse defeats it; soundness rests on the drop-order invariant. Corrected in `erased.rs` and in CLAUDE.md's unsafe-code entry, which carried the same overclaim. - `handles.rs`'s module doc still said tuple branches go through `build_branch` and panic on `ByteBounded`; they have used `build_branch_byte_aware` since the third round. - The same-item-type fan-out shapes (`OrderedBytesTuple2/3<OrdU64, ..>`) asserted only per-position downcasts and drain state, which an aliased branch handle satisfies. They now route a distinct value through each branch. - `maybe_instrumented` -- the constructor the branch builders actually call -- had no coverage on any of the three queue impls; one that dropped its `Some(metrics)` would have produced a silently unmonitored edge. - The edge-registry test asserted a count and `is_some()`, which swapped or duplicated consumer labels satisfy; it now pins both producer->consumer pairs. - `PipelineStats::record`/`record_error` index directly and panic on an out-of-range `StepIdx` while the detached recorders no-op; the asymmetry is deliberate and is now documented under `# Panics`. - `pool.rs`'s `sticky_exclusive_step` test helper took an `owner_idx` it discarded, with every call site passing 0. - `OrderedSource` / `OrderedPairSummer` declare `BranchOrdering::None` on an `OrderedBytesSingle` output, which maps to a direct branch -- so despite the shape's name these steps never engage the reorder stage. That is deliberate (the test reasons from FIFO order under byte backpressure) and is now stated, with pointers to the three places reorder restoration is actually covered. Sixth review round, four findings — the fused stall bound added last round, and three tests that asserted less than their messages claimed: - `fused.rs`: the stall budget was a hard-coded 5s, so a chain with a genuinely slow source had no recourse from a library-chosen policy. It now comes from `PipelineConfig::deadlock_timeout_secs`. `0` (the config default, meaning "monitor disarmed" on the scheduled path) selects a built-in default here rather than "unbounded", and the asymmetry is the point: the scheduled path has a monitor to arm and this driver does not, so a disableable bound would leave a fused wedge with nothing to detect it. Raising the number buys patience; there is deliberately no way to remove the bound. - `tests.rs`: the `build_two_input_handles` branch-mapping test asserted "branch 0 must feed input A and branch 1 input B" using `SumPairStep`, whose `a + b` is commutative -- `10 + 1` and `1 + 10` are both 11, so a swapped wiring passed while the message claimed the mapping was proven. Now uses a non-commutative `a * 1000 + b`; verified by swapping the wiring, which yields 1010 against the expected 10001. - `tests.rs`: `run_reports_finished_pipeline` sorted before comparing at every thread count, discarding a real guarantee at `threads == 1`, where `WideQueueSource` emits `n-1..0` through FIFO transports and one `ReportsFinishedBuffer` draining its `VecDeque` front-first. The single-worker case now asserts the exact descending sequence, so a reordering regression in the flush-first path is visible; above one worker the multiset assertion stands, because the Serial step's cross-worker re-dispatch interleaving is a valid scheduling detail. Stability confirmed over 40 consecutive runs. - `tests.rs`: `OrderedSink` recorded only `item.value`, leaving the out-ordinal sequence unasserted -- and the ordinal is precisely what `OrderedPairSummer`'s hold-and-retry path protects, since it assigns `next_out_ordinal` before the push and must reuse it on a retry. A regression that skipped or reassigned an ordinal left the value sequence intact and passed. It now records `(ordinal, value)` and asserts both. Seventh review round (a full re-review, not incremental), three findings: - `lib.rs`: the crate-root re-exports omitted `DetachedGroup`, which `Step::detached_group` returns while `Step` itself is flattened, and `OrderedBytesTuple2` / `OrderedBytesTuple3`, which sat beside an exported `Single` and `OrderedBytesSingle`. A step author implementing a `Detached` step or a 2-/3-way ordered fan-out had to mix flattened paths with `::step::` / `::outputs::` ones. All three added, pinned by a compile-time test that names every step-author type at the crate root — it fails to build if a `pub use` is dropped or a new output shape lands without one. A sweep for the same shape found two more unexported `pub` types, `queues::BoundedQueueHandle` and `reorder::ReorderCapHandle`, deliberately left alone: they are runtime budget plumbing named only through `runtime::contexts::RegisteredQueue`, which is not flattened either, and `ByteBoundedQueue` exposes inherent `limit_bytes` / `set_limit_bytes` / `current_bytes` so a direct user never needs the trait in scope. Promoting them would make the root surface less coherent, not more. The test records that decision so the next sweep does not re-litigate it. - `fused.rs`: `DEFAULT_STALL_BUDGET` raised from 5s to 60s. Because `deadlock_timeout_secs` defaults to 0, that default is what most fused runs actually get, and 5s is short enough for a source whose background reader blocks on a cold page cache or network-backed input to trip it — the idle timer only resets on progress, so one slow read is enough to fail a healthy run. The bound exists to catch a permanent wedge, not to police a slow source, and the costs are asymmetric: a too-tight budget loses output, while a generous one only delays the report of a wedge that has already hung. - `tests.rs`: `OrderedSink` carried two doc blocks with `#[derive(Clone)]` between them — the previous round's edit inserted the new doc after the attribute and left the stale `/// Records summed values.` line above it. Merged into one block before the attribute. A sweep for the same shape (a doc comment following an attribute) found no other site.
#714) Three additive entry points that `fgumi-sort`'s arena engine and the pipeline BGZF steps need, landed ahead of their consumers so the crate arrives complete rather than growing under each later port. `decompress_into_slice` is the fixed-slice analogue of `decompress_block_slice_into`: it decompresses a block straight into a caller-sized `&mut [u8]` instead of appending to a `Vec`, which lets the sort ingest path decompress directly into an arena slot rather than into a staging buffer it then copies out of. It keeps the integrity guarantee its siblings provide -- the payload must exactly fill `out` and match the footer CRC32 -- and takes the same deflate stored-block fast path for level-0 input. `uncompressed_size` is the accessor that makes that contract safe to satisfy. A caller has to size the slot before calling, and ISIZE is a `u32` sitting in the file, so sizing straight off the footer means allocating from unvalidated input -- up to 4 GiB from a corrupt block. This reads the footer off the same `&[u8]` (no owned `Vec`, so no copy of the block just to read four bytes) and bounds the claim to `MAX_UNCOMPRESSED_BLOCK_SIZE`. `decompress_into_slice` resolves the size through it too, so what a caller allocates and what the decompressor accepts cannot diverge. The stored fast path is now shared three ways. `copy_stored_and_verify`'s framing checks were inlined in its body, so the slice variant would have had to duplicate them; they move to `parse_stored_frame`. The BTYPE dispatch predicate becomes `is_stored_block`, and the inflate-then-verify tail becomes `deflate_into_slice_and_verify`, both shared with `decompress_and_verify`. Extracting only the framing check would have left the same drift exposure one level up. `InlineBgzfCompressor::recycle_buffer` lets a `take_blocks` consumer hand a drained block buffer back to the pool. Only `write_blocks_to` recycled before, so a consumer driving the compressor with `write_all` + `flush` + `take_blocks` left the pool permanently empty and allocated a fresh output `Vec` for every block. `write_blocks_to` now routes through the same method. The pool is bounded on both axes -- count and buffer capacity -- because bounding the count alone would let one oversized `Vec` handed in by a caller sit there for the compressor's lifetime. Rewriting `write_blocks_to` to call `recycle_buffer` (a `&mut self` method) means draining into a temporary rather than holding a borrow on `completed_blocks`. That makes the write-error path recoverable, so it is handled rather than left as-is: the failing block and its tail are restored to the queue instead of being dropped by the `Drain` guard. They are documented as safe to inspect, not to replay -- `write_all` can commit part of the failing block before erroring, and the whole block is re-queued, so writing the queue again would repeat those bytes inside a gzip member.
First slice of the arena sort engine (P3 of the feat-runall landing). `SortMergeSlot` is the bounded handoff between the phase-2 merge reader and its decompression workers: the reader pushes compressed spill blocks, workers claim and decompress them, and the slot carries the EOF and error signalling that tells a drained consumer whether the stream finished or failed. It lands ahead of the engine that uses it, alone. The other eight modules in this phase reach into the `external.rs`, `worker_pool.rs`, `segmented_buf.rs`, and `inline.rs` rewrites, so they cannot compile against the current tree and have to arrive with those. This one has no such dependency -- it needs only `codec` -- which makes it the one piece of the engine that can be reviewed on its own rather than inside an 8,500-line diff. That matters more here than elsewhere: it is lock-free concurrent code whose bugs are interleavings, not lines. It arrives with its model check. `tests/loom_merge_slots.rs` drives the real `SortMergeSlot` under `loom`, which enumerates thread interleavings systematically rather than leaving them to timing luck. `loom` is a `cfg(loom)`-only dependency so the normal build never sees it, and the test target is empty without the flag -- which is exactly why `check.yml` now runs it: a broken model, or a `cfg(loom)` build that stopped compiling, is otherwise invisible, since the target compiles to an empty binary and reports success. Loom explores a bounded state space deterministically -- four of the five models under a preemption bound -- so the same input yields the same verdict, and unlike `miri.yml` and `stress.yml` it can gate a PR instead of running on a schedule. Changes made to the module relative to the source branch, all of which tighten rather than alter behaviour: - `bp_insert_drain_finalize` documented that `blocks.len()` need not equal `count`. Delivering fewer blocks than reserved releases the reservation for sequences never inserted, so the slot finalizes a *clean* `queue_eof` over the gap and a truncated spill file is indistinguishable from a complete one. Releasing more wraps `in_flight` to `usize::MAX`, after which the finalize predicate can never hold and the merge polls forever. Both are now documented as contract violations and caught by assertions that are ALWAYS ON, release included: the failure they catch is silent record loss or a wedged merge, neither has an in-band signal, and release is what ships. Two compares per batch (not per record), both before any lock is taken, so a violation cannot poison a mutex. The `#[should_panic]` tests lost their `#[cfg(debug_assertions)]` and now run in release too -- where they previously did not exist at all. - `bp_add_in_flight` and `bp_set_reader_eof` were `pub` despite their own docs saying the call order is safety-critical and `bp_commit_read` existing to encapsulate it. They are `pub(crate)` now; nothing outside the module calls them on any branch. - `loom_decomp_error_beats_clean_eof` asserted `!(drained && !errored)`, which is unsatisfiable either way and stayed green when the `decomp_error` guard was deleted from `is_drained`. It now asserts `!drained` directly -- the producer's only path is the error path, so no interleaving may report a clean drain -- and fails under that mutation. A reversed-read-order consumer models the documented order-independence. - The module header claimed this is the live production Phase-2 and that the reorder buffer and in-flight counter were eliminated. Neither holds here: the module has no callers, `worker_pool::Phase2FileState` is the live path, and the struct declares both fields for the block-parallel path. Header rewritten to describe this tree. - `is_drained`'s doc claimed both flags are read under the `decompressed` mutex; `queue_eof` is read on the unlocked fast path. Documented why that is sound instead of asserting something untrue. - `PHASE2_DECOMP_CAP` is 32 here and 8 in `worker_pool.rs`, and the crate root exports this one while the other governs the live merge. The export stays -- `fgumi-pipeline-io` imports it in a later phase -- with a doc note on which is in force. - `decomp_error_flag_set_and_read` tested `AtomicBool`, not the slot; it now exercises `has_error`. An assertion in `bp_reorder_window_is_bounded_under_straggler` was vacuous because `reader_eof` was never set; the test now sets it and checks that EOF does not finalize over an in-flight straggler, then that every block reaches the FIFO once it lands. `fifo_order` pushed into and popped from `decompressed` by hand, asserting that `VecDeque` is a queue rather than anything about the slot; it now drives out-of-order completions through the reorder buffer. Both of those closed on a block COUNT while every block was identical bytes, so a drain that dropped one block and duplicated another, or delivered the window backwards, passed either way. Blocks are stamped with their sequence number and the assertions are on identity and order: reversing the drain fails both now and failed neither before. - `SortMergeReader` is reachable through a public field but was not re-exported, so callers could hold the value without naming the type. - The loom test's header said a protocol that failed to finalize `queue_eof` "surfaces as the `queue_eof` assertion below, not as a hang". It cannot: `consume_until_drained` breaks only on `is_drained()`, and `run_model` joins the consumer BEFORE its assertions, so a non-finalizing protocol spins the consumer, never joins, and never reaches them. The real symptom is a loom `max_branches` panic naming the harness rather than the protocol. Both that note and the same claim on `consume_until_drained` itself now say so. - `drain_locked_and_finalize` stops once the FIFO reaches `PHASE2_DECOMP_CAP`, so a slot can be at reader-EOF with `in_flight` zero and still owe blocks that did not fit. Finalizing there would strand them; completion instead depends on the driver draining again once the consumer makes room. No test reached it -- the cap is 32 and every test and model used at most five blocks -- so it is now covered at both boundaries (FIFO exactly at cap, and one short of it), and each case fails if the `reorder`-empty guard is removed. - `docs/design/sort-phase2-unification-deferral.md` is added: the module header cites it twice, and the source branch pins it in `.gitignore` precisely to stop that citation dangling. Its two claims that do not hold in this tree are scoped rather than ported verbatim: it lists a streaming driver whose `external.rs::MergeDriver` half does not exist on this branch (so a reader could apply its deadlock analysis to the uncalled `SortMergeSlot`), and it called the pool path "immune" to deadlock where what was shown is that one analyzed failure mode does not arise.
Second slice of the arena sort engine, and the last one that can land before the engine itself. Two pieces, both leaves: the arena pool plus the `SegmentedBuf` API it needs. `arena_pool.rs` bounds and reuses the Phase-1 sort arenas. The spill path fills a multi-GB `SegmentedBuf`, hands it downstream as an `Arc`, and previously minted a fresh buffer for the next fill via `mem::take`. With block-parallel spill running async that leaves two full arenas live at the Phase-1/Phase-2 boundary -- the in-flight chunk, pinned by slow compression through the gather's bounded queue, plus the next fill -- which is where the ~2x peak RSS came from. The pool caps live arenas at `capacity` (default 1, matching the legacy one-arena-at-a-time model) and recycles their storage: `try_acquire` returns a reset-for-reuse buffer or `None`, and the caller backpressures until an in-flight chunk's `Arc` drops and `PooledSegmentedBuf::drop` returns its arena. `segmented_buf.rs` gains the two methods those depend on -- `reserve_full_capacity` and `reset_for_reuse` -- and nothing is removed or changed: the diff is +475/-22 and the 22 are the doc rewrite around the new API. Its existing callers (`inline.rs`, the spill path) are untouched, and the crate's 7,836 tests pass either way. Changes relative to the source branch: - `PooledSegmentedBuf::drop` nested `if let Some(buf) = ...` inside `if let Some(pool) = ...`, which `clippy::collapsible_if` rejects under this tree's toolchain. Rewritten as a let-chain rather than allowed, per CLAUDE.md's rule to adopt the idiom the newer compiler unlocks. The `take()` still runs unconditionally, so an unpooled buffer drops exactly as before. - Module doc comments referred to "increment 1/2/3" and "lever-1" -- the source branch's internal campaign numbering, which names nothing a reader of this history can look up. Reworded to describe the work. - `arena_pool` was `pub` while `segmented_buf` stayed `pub(crate)`, so the pool's public signatures (`try_acquire`, `pooled`/`unpooled`, `Deref::Target`) named a crate-private type and its doc links pointed at one. CI's `docs` job sets `RUSTDOCFLAGS=-D warnings`, which turns that into a build failure -- a plain `cargo ci-doc` only warns, which is why it was missed locally. `segmented_buf` is promoted to `pub` with `SegmentedBuf` re-exported, matching what the engine does; demoting `arena_pool` instead would have to be undone one commit later. - CLAUDE.md gains the `segmented_buf.rs` unsafe allowlist entry for `grow_uninit` and `slice_mut`, which this commit introduces, plus a note that the counts given are production sites and a raw grep also sees the test-only ones. - A comment claimed `slice_mut_concurrent_disjoint_writes_are_sound` is race-checked under miri. It is not: `miri.yml` covers `fgumi-raw-bam` and `fgumi-pipeline-core` only. Corrected to say what the test does and does not establish. A deep review pass found three things worth recording. `ArenaPool::try_acquire` returned a bare `SegmentedBuf`. `made` is only ever incremented and the sole path back to the free-list is `PooledSegmentedBuf::drop`, so any drop that bypassed the wrapper -- a `?` on a malformed record, an early return, an unwind -- retired that slot permanently. At the default capacity 1 the pool is then empty forever, and since the documented response to `None` is to backpressure until an arena returns, the caller spins rather than fails: a hang, not an error. The engine's own call site had exactly this shape. `try_acquire` now returns the RAII wrapper and `PooledSegmentedBuf` gained `DerefMut` so the fill happens through it; the leak is no longer expressible. `slice_mut`'s `# Safety` listed in-bounds-ness and range disjointness. It also requires that no `&mut self` method run concurrently: `grow_uninit` can reach `advance_segment`, which pushes to `self.segments` and so reallocates the OUTER `Vec<Vec<u8>>` a concurrent reader is indexing -- a use-after-free, not a stale read -- and its `set_len` races that reader's bounds check. Neither is fixed by `reserve_full_capacity`, which only pre-sizes the inner buffer, yet `grow_uninit`'s pointer-stability note claimed it was sufficient for the concurrent case. Miri reports UB on that pattern with both callers honouring every documented precondition. The contract now states the third requirement and describes the shape that is actually supported (reserve every slot for a segment, then hand them out); CLAUDE.md's allowlist entry is corrected to match. `block_offsets.rs` is deleted rather than landed. It re-derived `SegmentedBuf`'s offset arithmetic in a second place, and its seal test claimed parity with the `memory_usage() >= memory_limit` spill trigger while counting only arena data bytes -- `memory_usage` also counts the per-record ref vector, so every planned arena would overshoot the budget the sibling `ArenaPool` exists to hold. It had no caller in this commit, none in the arena engine, and none on the source branch, so the parity bug was unreachable and the correct fix unknowable; `make_room` stays the one place that arithmetic lives. Test gaps the same pass found, each verified by mutation: - `clear` truncates to one segment, so it must also rewind the write cursor or the next write indexes out of bounds. Only tested on a single-segment buffer, where `cur` is already 0 -- deleting the reset passed. Now covered across several segments, and the new test is the only thing that fails under that mutation. - `reserve_full_capacity` reserves `segment_size - seg.len()`, and `reserve_exact` takes an amount ADDITIONAL to the length, so dropping the subtraction over-allocates a partly-filled segment by up to 2x. Only the empty case was tested, which cannot see it. - `capacity` was only ever exercised at 1, where "hand out one arena" and "respect the capacity" are indistinguishable; an implementation ignoring the field passed. Now a case table at 2 and 5. - `slice_mut_spanning_segment_boundary_panics` asserted only `result.is_err()`, which any panic satisfies, and its scenario over-read within one segment rather than crossing a boundary. Both cases now assert on the guard's own message, and one genuinely crosses. Lower-severity findings from the same pass, and a class sweep over what they had in common: - `PooledSegmentedBuf::pooled` was `pub`, so a caller could wrap a buffer this pool never handed out; `release` pushes it onto the free-list and `try_acquire` pops from `free` BEFORE consulting `made`, which raises the live-arena count above `capacity` and defeats the RSS bound. It is private now -- `try_acquire` is the only way to get a pooled wrapper, so every buffer on the free-list came from this pool. - The zero-length early return in `slice`/`slice_mut` skipped all bounds checking, so `slice(1_000_000, 0)` silently returned an empty slice where it used to panic. Narrowed to the case it exists for (the one-past-the-end offset a zero-length reservation produces). - `advance_segment`'s "landed on a non-empty segment" check was debug-only. A non-empty segment breaks the `cur`/`total_len` relation that `make_room` maintains, so `locate` resolves every later offset into the wrong segment and the sort emits wrong bytes with no error. It runs once per segment transition -- once per 256 MiB on the sort path -- so it is always on now. - The `cur` field claimed `cur == total_len / segment_size` "always holds". It holds at segment transitions, which is what keeps `locate` exact, but it is not a standing invariant. Reworded to say which. - The module header said appending "never moves existing bytes". True of earlier segments, which is the point of the design; not true within the current segment, whose `Vec` still reallocates unless pre-sized. - `ArenaPool::new` clamped `capacity` but not `segment_size`, though `SegmentedBuf::with_capacity` clamps it downstream anyway. - `PooledSegmentedBuf`'s doc named `from_owned_records`/`empty` constructors that do not exist. - `Deref`/`DerefMut` are the only way a consumer reaches a pooled arena and were never exercised; the error-path test asserted only `is_err()`, which "pool empty" also satisfies. Both are covered and discriminating now. - CLAUDE.md's allowlist cited `≈line` pointers that had drifted ~50-60 lines. They name the function instead: an approximate line number goes stale the moment anything above it moves, and a confidently wrong pointer is worse than none. The `fgumi-raw-bam` entry counted four sites without saying two are tests, which contradicted the production-only counting rule; it now says which. - `release` propagated a poisoned pool mutex, and it runs only from `PooledSegmentedBuf::drop`. An arena dropped while a panic is already unwinding would panic a second time and abort the process outright. The poisoning is reachable rather than theoretical: `try_acquire` allocates the arena *inside* the critical section, and `Vec::with_capacity` panics on capacity overflow, so a large enough `segment_size` poisons the lock. `release` recovers the guard instead -- sound because the protected state cannot be torn (a free-list of reset buffers plus a counter). Its only other lock, `try_acquire`, keeps `expect`: it is not on a drop path, and a poisoned pool should surface there. `free_len` recovers too, so a test can still observe the pool afterwards.
The engine itself -- the last and largest slice of P3. Six modules plus the `external.rs`, `worker_pool.rs`, `inline.rs`, and `keys.rs` rewrites that carry them. `SortMergeSlot` and its loom model (#718) and the arena pool and block-offset planner (#719) came out of the same upstream commit and were split off ahead of this one; this is its remainder. `chunk_sorter.rs`, `ref_sort.rs`, `template_arena.rs`, `spill_block.rs`, `spill_block_reader.rs`, and `sync_spill_writer.rs` could not land separately: the first three form a dependency cycle among themselves and all reach into the rewritten `external.rs`, the last three need `external.rs` and `worker_pool.rs`, and `external.rs` in turn depends on the new modules. This is a three-way merge, not a port. The source commit predates 30 commits of `fgumi-sort` work on main -- 18 on `external.rs` alone -- so wholesale copying would have silently reverted them. Where both sides moved: - `keys.rs` takes both improvements to the queryname key: upstream's inline `SmallVec` name buffer and main's `pos: u32` ingest position (#621), which is what makes name+flags a total order and lets the chunk sort be unstable. Taking either side alone drops an optimization or reintroduces nondeterminism. The padding test asserted 32 bytes on the reasoning that a `Vec<u8>` is 24; the key is larger, but the test's claim -- that `pos` is free -- still holds because name+flags alone already round up. The expectation and its derivation are corrected and the padding pinned directly. - `header_ss_tag` is byte-identical to main's, so the `@HD SS` spelling fixes (#514, #567) survive. Verified by diffing against the base. - `inline.rs` drops the `RecordBuffer::max_sort_key` reset: both sides optimized the same thing and main's derives the radix bound inside the first counting pass (#622), needing no stored field. - `lib.rs` re-exports are a union. The engine's list dropped `format_thread_counts` (#691), `SortMergeReader`, and the `reader.rs` entry points, all of which the root crate imports. - `worker_pool.rs::set_main_thread` is dropped rather than ported. It needs an `ArcSwap` field type from an upstream change main never took and has no callers; its own doc says its production caller was retired. `smallvec` is a crate-local dependency, not workspace-inherited -- `keys.rs` is its only user -- with the `const_generics` feature, since the default `Array` impls cover 1..=32 then 36, 64, and the inline name capacity is 44. That capacity fits 100% of read names in every local benchmark dataset; on native Illumina names (37-41 bytes) it is 9.5% less CPU and 4.4% less peak RSS end-to-end, against ~4% more RSS on short SRA-normalized names. Doc claims that were true when written and are not now: `merge_slots.rs` said `SortMergeSlot` has no callers (it has two, both test-reachable); `worker_pool.rs` said the pool path no longer drives any production sort (here it is the only thing that does, until the typed-step `SortMerge` consumer lands with `fgumi-pipeline-io`); and the Phase-2 unification deferral doc said `external.rs::MergeDriver` does not exist. CLAUDE.md gains the `ref_sort.rs` and `segmented_buf.rs` unsafe allowlist entries, plus a note reconciling a raw grep's 22 `#[allow(unsafe_code)]` sites against the 7 production sites the allowlist names -- the difference is entirely tests exercising already-approved APIs.
…mulated (#722) `cargo ci-doc` builds the workspace docs with `RUSTDOCFLAGS="-D warnings"`, which reads as "every rustdoc warning is an error." It was not. Rustdoc only checks intra-doc links on items it actually renders, and without `--document-private-items` that means public items only -- so a broken link in the doc comment of a private fn, a private field, or a `#[cfg(test)]` helper never reached the linter. Thirty such links had accumulated across the workspace, all of them invisible to a green docs job. This adds the flag to the `ci-doc` alias and fixes every link it surfaces, so the job now checks what its name implies. This is a gap in the check, not a pile of typos. It was found only because a broken link in `crates/fgumi-sort/src/inline.rs` was caught by review on #720 -- CI could not have caught it, and could not catch the other twenty-nine either. Closing the gap is the point; the fixes are the cost of closing it. None of the thirty were reachable defects: the cost was a doc build that quietly under-checked itself, plus links that render as literal brackets instead of hyperlinks. Four kinds: - Not in scope at the link site, so the path was required -- `PipelineError::*`, `EdgeMetrics`, `MAX_ARITY`, `TwoInputHandles`, `SortWorkerPool::phase1_queue_depths`, `PooledBamWriter::new_indexing`. - Associated fns linked bare instead of through `Self::` -- `merge_chunks_generic`, `merge_chunks_with_index`. - A symbol in another crate -- `read_raw_record` lives in `fgumi-raw-bam`, not `fgumi-sort`. - Prose that rustdoc parses as a reference-style link -- byte-layout notation `[len(4)][record(len)]`, index expressions `buffer[0]` and `R1[i]`, and ``[`X`](path)`` pairs whose explicit target is redundant because the link text already resolves. Touches `fgumi-sort`, `fgumi-pipeline-core`, `fgumi-bam-io`, and the main crate. Docs and one cargo alias only -- no functional change. Based on `main-runall` rather than `main` because eight of the fixes are in `fgumi-pipeline-core`, which does not exist on `main`. Three further instances of the same class, in files that #720 introduces, are fixed in #720 itself rather than here.
* feat(pipeline-io): add the fgumi-pipeline-io crate
Adds the BAM-pipeline I/O layer for the typed-step pipeline: the source,
sink, and sort building blocks plus the record-batch and BGZF-block buffer
types exchanged between steps.
* source — BAM ingest (ReadBgzfBlocks, read_bam*) turning a reader into
decompressed blocks.
* sink — BAM output (WriteBgzfFile).
* sort — the in-pipeline sort steps (SortBuffer, CompressSpill,
SortSpillDecompress, SortMerge) over the arena ingest, spill-write,
spill-gather and boundary-finding machinery.
* types — RecordBatch, DecompressedBlock, BgzfBlock and friends.
Nothing constructs these steps yet: `fgumi sort` still drives the legacy
file-to-file sorter, and routing the command through the arena chain is
root-crate wiring that lands separately. This commit adds the machinery
only, so the crate is exercised by its own tests rather than in production.
Ported from origin/feat-runall f50df03, but NOT byte-identical to it. The
divergences are listed below so a future port from that branch knows where
to expect conflicts.
Toolchain / workspace drift (no behaviour change):
* ArenaPool::try_acquire already returns the pooled wrapper on this base,
so the manual PooledSegmentedBuf::pooled wrap is dropped.
* Three collapsible_if sites rewritten as let-chains rather than
suppressed, per the repo's rule about adopting newer-compiler idioms.
* Three rustdoc intra-doc links fixed. One was a genuine doc inaccuracy:
FindBoundariesAndSort's doc described the wrong sort entry point.
Behaviour changes made under review, which feat-runall does not have:
* TemplateChunks::push returned unreachable!() on a template lane-variant
change. The --key-types variant is global to a sort, so a change means
phase 1 and the merge disagree about key width and the merge would emit
silently mis-ordered output. It now returns io::Result with an
InvalidData error naming both variants, threaded through
MemoryChunksByKind::push and absorb_phase2_event — the same fail-closed
treatment ensure_single_lane and build_driver already had.
* SortSpillDecompress::emptiest_first_order now filters slots that have
reached queue_eof before reading FIFO lengths. The registry is
append-only, so a drained slot otherwise stays in the scan for the rest
of the run, taking its FIFO lock on every dispatch only to be rejected.
Scheduling-only: an EOF slot cannot progress, so no output changes.
* ReadBgzfBlocks held its reader and finished flag in an Arc<Mutex<..>>
and an AtomicBool. It is a Serial step, so the runtime drives one shared
instance and never clones it; both are now plain fields, the misleading
"was clone() called?" panic message is gone, and the reader is released
at end of stream instead of held for the whole run.
* The residual-emission rule is split out of ChunkSorter::take_residual_chunks
into a free function so it can be unit-tested without building an
accumulator. Same rule, same behaviour.
Also registers the crate in scripts/publish-crates.sh, after fgumi-sort and
before fgumi-consensus. The list is topologically ordered and CI's
publish-order job fails when a workspace crate is missing from it; the new
crate depends on fgumi-bam-io, fgumi-bgzf, fgumi-pipeline-core,
fgumi-raw-bam, and fgumi-sort, so that is the earliest valid position.
* chore(sort): drop stale allow(dead_code) suppressions
These six attributes predate the pipeline-io machinery and no longer
suppress anything: clippy over --all-features --all-targets is clean
without them. Removing them restores dead-code coverage for
SegmentedBuf's inherent impl, TmpDirAllocator::with_recheck_interval and
mark_full, PrefetchReader::bytes_consumed and consumer_stalls, and the
memory-debug-gated print_mi_stats.
* test(pipeline-io): cover record framing, step wiring, and spill plumbing
The ported crate arrived at 80% patch coverage, below the 90% gate, and the
gap was concentrated in the code that decides whether sorted output is
byte-identical.
boundaries.rs was the worst of it: 31.6% of lines and only 22% of functions,
with a single test covering `bam_header_len` and nothing at all for the
scanner that decides where records begin and end. It now sits at ~99% lines
and 100% functions, with cases for header skipping, cross-block record
carryover, trailing partial records, incomplete headers, and every `finish`
rejection path. The existing sequential-assert test was refactored into an
rstest case table so a failure names the scenario.
Also covered: `ReadBgzfBlocks` driven through a real pipeline (its `try_run`
had no coverage — the previous test hand-rolled an equivalent read loop
instead, so it validated `read_raw_blocks` rather than the step, and has been
removed), the SpillWrite / SpillGather / SpillBlockCompress / ReadBlocks /
InflateToArena step wiring, `SpillWrite::open_file`'s refusal to reuse a
file_id, `SpillGather`'s dense ordinal minting and empty-chunk skip, the
TemplateCoordinate framing path (the one variant that does not share
frame_keyed_record_into), `SpillBlockEvent`'s heap_size and its Ordered
delegation (which guards an infinite-recursion trap), the residual-emission
rule including its template-coordinate exception, the block_batch clamp, the
shared-Arc fail-closed guard, and the over-reserve guard on the k-way merge
path rather than only the fast path.
Two fixes to the tests themselves, both of which made existing assertions
weaker than they looked:
* synthesize_sized_records built every record with tid = -1, so
extract_coordinate_key_inline returned RawCoordinateKey::unmapped()
(u64::MAX) for all of them, `pos` was never compared, and the coordinate
parity cases collapsed into one equal-key bucket — they would have passed
with a broken key comparison. Records are now mapped across four
references, with every eighth left unmapped to keep tie coverage.
* The queryname branch only length-checked the RecordBatch output, so a
framing divergence between SortMerge<BlockOutput> and
SortMerge<RecordBatchOutput> that preserved record count would pass. It
now compares the two framings as multisets.
The multi-minute soak, matrix, and proptest suites move behind the crate's
new `stress-tests` feature, matching the root crate's existing convention:
the default test target drops from ~30s to ~2s for this crate, and
`cargo ci-test-stress` still runs them. The coverage job explicitly enables
that feature — those suites reach ~200 lines nothing else covers, so
measuring without them under-reports coverage for code that is tested, just
not on the PR-latency path.
Where a StepCtx is required the crate's existing convention is followed:
test the StepCtx-free core, or drive the step through Pipeline::builder.
All test data is constructed programmatically; no fixtures are committed.
`merge_phases.rs` (#706) opens with an intra-doc link to `crate::external::SortPhaseTimer`, but the struct is private to `external`, so the path is not nameable from another module and rustdoc cannot resolve it. Widen it to `pub(crate)` — the narrowest visibility that makes the existing cross-reference valid. No public API change. This surfaced only when #706 was rebased onto this branch: `main`'s `ci-doc` alias does not pass `--document-private-items`, so the broken link is invisible there, while this branch has checked private-item links since #722. The same gap will keep producing this class upstream until `main` adopts the flag.
…#734) The typed-step pipeline tree needs `DecodedRecord`, `GroupKey`, `GroupKeyConfig`, `Grouper`, `compute_group_key_from_raw` and `name_hash_key`. Today those live in `src/lib/unified_pipeline/{base,bam}.rs` (and `name_hash_key` does not exist at all), so the new tree would have to depend on the tree it is meant to replace. Move them next to the raw-record helpers they operate on, in `fgumi-bam-io`, along with the `LibraryLookup`/`LibraryIndex` read-group-to-library machinery the group key depends on. Both files are self-contained — their only imports are `noodles::sam` and each other. `unified_pipeline` keeps its own copies for now; the duplicate is retired when that tree is deleted.
…l steps (R1a) (#735) * feat(template): expose heap_size as the single byte-budget source The typed-step pipeline budgets its queues on real allocated bytes, so `BamTemplateBatch` needs a `Template` heap footprint it can sum. The capacity-based computation already existed, but only inside the `MemoryEstimate` impl, where the pipeline cannot reach it. Lift it to an inherent `heap_size` and make `estimate_heap_size` delegate, so the queue budget and `MemoryEstimate` cannot drift apart. Adds the agreement test, which also re-derives the expected floor from the components rather than from `estimate_heap_size` (that would be tautological once the trait delegates). * feat(pipeline): add the typed-step pipeline tree with its foundational steps First slice of the umbrella-side port. Additive: nothing in `src/lib/commands` routes through this tree yet, so every command still runs on `unified_pipeline`. `pipeline/mod.rs` re-exports `fgumi_pipeline_core as core`, which is what lets the step sources compile unmodified — their `crate::pipeline::core::…` paths resolve through the shim rather than needing a rewrite. Ports `steps/{types,tuning}` plus the `sink/`, `sort/`, `bgzf/`, `boundaries/` and `parse/` subtrees. `source/`, `group/`, `correct/` and the closure-driven mid-steps follow; `align_and_merge` waits on `crate::aligner`. Two deliberate deltas from the source branch: - `boundaries/bam.rs` uses a let-chain instead of nested `if`s. The source predates the toolchain that unlocked them, and `ci-lint` rejects the nested form as `collapsible_if`; per CLAUDE.md we adopt the idiom rather than suppress the lint. - `tuning.rs` demotes the `PipelineConfig::auto_tuned` intra-doc link to a code span. The referent lives in `unified_pipeline`, which this port exists to retire, so linking it would only break again on deletion.
* feat(fastq): budget FASTQ records and templates by heap size The typed-step pipeline byte-bounds its queues on `HeapSize`, which `FastqRecord` and `FastqTemplate` do not implement — only `MemoryEstimate`, which the framework cannot see. Add `HeapSize` impls that delegate to `MemoryEstimate::estimate_heap_size` rather than restating the computation, so the queue budget and the memory estimate stay a single source of truth. Mirrors the treatment `Template` already received. * feat(pipeline): port the read-side source steps Adds `steps/source/`: the FASTQ readers (`ReadFastqInputs`), the serial chunk-level pairing step (`PairRawFastq`), the parallel parse+zip step (`ParseAndZipFastq`), SAM chunk reading, and a re-export shim for the BAM source that already lives in `fgumi-pipeline-io`. Still additive — no command routes through this tree. Deltas from the source branch, all forced by this branch's stricter gate: - Four nested `if`s in `pair_fastq.rs` become let-chains. `ci-lint` rejects the nested form as `collapsible_if`; per CLAUDE.md we adopt the idiom rather than suppress the lint. - Three module-level intra-doc links (`ReadFastqInputs`, `ParseAndZipFastq`, `PairRawFastq`) gain reference-style definitions. They resolved only inside the item docs that already defined them; the module docs had none, which `ci-doc` catches here because this branch checks private-item links.
…one definition (#741) `DecodedRecord`, `GroupKey`, `GroupKeyConfig`, `LibraryIndex`, `LibraryLookup`, `build_library_lookup`, and the `Grouper` trait each existed twice: once in `fgumi-bam-io` (ported in #734) and once in `unified_pipeline` / `read_info`. The copies were textually identical apart from field visibility and doc formatting, and the duplication was deliberate — #734 left the umbrella copies in place so nothing had to be rewired. That worked while the ported step library only ever *produced* these values. It stops working at the first ported step that hands one to an umbrella grouper: `GroupBam` calls `TemplateGrouper::add_records`, and the two identical definitions are two incompatible types, so the call cannot be written at all. Keep the `fgumi-bam-io` definition and re-export it from both former homes, so every existing `crate::unified_pipeline::DecodedRecord` and `crate::read_info::LibraryIndex` path resolves unchanged. Two supporting changes fall out: - `GroupKey::has_mate_position` only existed on the umbrella copy; ported to `fgumi-bam-io` so the surviving definition is a superset of both. - `grouper.rs` reached into `DecodedRecord`'s fields directly, which are private on the `fgumi-bam-io` copy (made so during #734's review). Those five sites now go through `into_raw_bytes` / `record` / `raw_bytes`, which is what the accessors are for. Net -483 lines with no behavior change; the full suite is unchanged at 8322.
The chain builder needs each stage's tuning knobs without depending on how they were supplied, so it can construct a chain with no CLI struct behind it. This adds one plain `<Command>Options` struct per stage, plus a `to_<command>_options()` method that projects the parsed flags into it, for correct, group, simplex, duplex, codec, filter, and zipper. The structs deliberately do not derive `clap::Args`. Flattening them into the command structs would move the fields off those structs and rewrite every `self.<field>` reference and its tests, which buys the chain builder nothing — it only ever reads the values. That refactor belongs with the fused command that actually needs prefixed `--<stage>::<flag>` flags. As a result this change is additive: no CLI surface moves, and the existing suite covers the commands unchanged. Two option structs hold resolved rather than raw values, because grouping cannot be configured from the raw flags alone: - `GroupOptions::min_map_q` applies the default `execute` would apply. - `GroupOptions::effective_strategy` / `effective_edits` carry the `--no-umi` and identity-implies-zero-edits rules. Both now come from `GroupReadsByUmi::resolved_min_map_q` and `resolve_strategy_and_edits`, extracted from `execute` so the command and the chain builder read them from one place instead of duplicating the rules. `execute` keeps its own validation and override logging; the new methods only compute. A case table pins both sets of rules. `GroupOptions::index_threshold` keeps the richer `IndexThreshold` type rather than a bare count: the flag also accepts `always` / `never`, which a number cannot express. Each projection is covered by a test that drives the command through `try_parse_from` rather than a struct literal, so a renamed or unwired flag fails the test instead of silently compiling. `Strategy` gains `PartialEq`/`Eq` so the resolution table can assert on it.
…g at teardown (#745) The scheduled runtime's deadlock monitor shipped disarmed, so a wedged pipeline hung silently and forever. The fused path, which has no monitor, carries an unremovable 60s stall bound for exactly that reason — the two runtimes disagreed about whether a wedge should be survivable. Two costs kept the monitor off, and this removes both. Liveness came from `PipelineStats`, which is deliberately not free: `dispatch_one_step` times every dispatch with `Instant::now()` and gates that on `stats.is_some()` to keep the uninstrumented path zero-cost. So arming the monitor meant buying per-dispatch timing. The monitor only ever needed one question answered — has anything progressed? — so it now reads a dedicated `LivenessCounter`: one counter per worker, each padded to its own cache line so a bump is an uncontended increment, summed by the monitor once per poll. Instrumentation stays opt-in. Teardown polled a stop flag in 25ms slices, so `Pipeline::run` waited up to a full slice for the monitor to notice before joining it — dead time on the critical path of every run. A condvar wakes it immediately. Both measured on a quiet c8g.4xlarge (aarch64) with a dispatch-overhead benchmark added here, each configuration run three times against a shared baseline, with a clean-vs-clean control first to establish the noise floor: liveness counter alone no effect distinguishable from noise monitor armed, polled +80% (4 threads), +23% (8 threads) monitor armed, condvar within noise The +80%/+23% figures reproduced to three significant figures across reps; the counter's did not, which is what separates signal from drift here. The default stays 0. What now blocks arming it is neither cost: `MonitorBlindTransport` makes an armed monitor reject any chain with a non-`ByteBounded` edge, because its stall verdict counts bytes in flight to tell idle from wedged. `Process2` uses `CountBounded`, so defaulting it on would fail those chains rather than watch them. That is its own change. The benchmark is new because the pipeline had none, and no command drives it on this branch yet, so a hot-path change could not otherwise be measured.
Adds the typed steps that sit between the read side (ported in the previous change) and the write side: template grouping, the closure-driven `Process*` family, and the serialization steps that turn grouped templates back into records. - `steps/group/` — `GroupBam`, `GroupByPosition`, `GroupByQueryname`, and `GroupByMi` - `steps/process.rs` — `Process`, `ProcessOrdered`, `ProcessWithWorkerState`, and the 2-/3-input variants - `steps/serialize.rs`, `steps/serialize_processed.rs`, `steps/templates_to_records.rs` - `steps/tests.rs` — full-chain tests over the newly wired steps Scope is deliberately limited to the steps the chain builder actually constructs, so two upstream files are left out: `coalesce.rs`, a `CoalesceBytes` byte-batching step that has no caller anywhere, and `roundtrip.rs`, whose only consumer is the `compare bam-roundtrip` command and which therefore travels with the command rewiring rather than with the steps. `steps/mod.rs` records both reasons. Two supporting changes come along because the new steps need them: - `fgumi-raw-bam` gains `find_two_string_tags_in_record`, a single pass that locates two string tags at once, plus its tests. - `mi_group.rs` gains `MiKey`/`MiTransform`. `MiKey` compares the raw tag bytes and only converts lossily to build the display label. The previous shape converted before comparing, so two distinct MI tags that are both invalid UTF-8 collapsed to U+FFFD and grouped together. The nested `if` blocks in the ported files are written as let-chains, which this toolchain's `clippy::collapsible_if` requires.
The type was renamed CodecStats -> CodecConsensusStats on this branch, but the doc link in the per-molecule tally struct still referenced the old name.
Adds `pipeline/steps/correct/`, the typed `Step` wrapping UMI correction, completing the mid-step group #743 left unfinished. `steps/mod.rs` recorded `correct/` as blocked on the command-layer options refactor; `CorrectOptions` arrived with #744, so the block is lifted. The step needs a little more of `commands::correct` than that type alone. `feat-runall` refactored that file and `main` has since developed it independently -- main's copy is the larger of the two -- so taking the upstream version wholesale would revert main's work. Widen only what the step names instead: - `pub(crate)` on `credit_umi_metrics`, so the step and the legacy `execute` path share ONE definition of fgbio's per-segment UMI accounting rather than each carrying its own. - `pub(crate)` on `RejectionReason`, `TemplateCorrection` (plus its `matched`, `matches` and `rejection_reason` fields), and the `CollectedCorrectMetrics` fields, plus the three associated functions the step calls. - `CollectedCorrectMetrics::merge_into`, to fold the per-thread accumulator slots once the pipeline has drained. This is genuinely new -- the legacy path aggregates its slots inline. The legacy path is unchanged. Metrics follow `execute` exactly, which follows fgbio: every template counts once in `templates_processed`; per-UMI crediting happens for every template that reached matching, *before* the keep/reject decision, so a template rejected on one segment still credits its matched segments and credits the all-`N` bucket only for the segments that actually failed; and missing-UMI and wrong-length templates credit no per-UMI bucket at all (`CorrectUmis.scala:199-202`). `per_umi_crediting_matches_fgbio` pins each of these; its `AAAA-TTTT` case is the discriminating one, since a template rejected on its second segment must still credit `AAAA` for its first. Both output shapes share one `run_batch`, which takes an optional rejects sink, so that accounting has a single definition and cannot drift between them. Forward-ports one behaviour the ported source predates: `apply_correction_to_raw` grew an `original_tag` parameter, so `CorrectStepConfig` carries `original_tag` beside `umi_tag`, both derived from `Target`. The module is `#![allow(dead_code)]` until the command rewiring gives it a caller: the surface cannot be `pub` without leaking `pub(crate)` command internals. The tests build `CorrectOptions` explicitly rather than via `Default`, which upstream got from the `multi_options` macro -- deriving it here would yield `cache_size: 0`, contradicting the flag's `default_value` and tripping a ported assertion that the default is non-zero. No command is rewired, so nothing executes on the ported path yet; that starts at R2.
The last_emitted_at field's doc linked [`Self::finish`], but Self is LadderLibraryState, which has no finish method. Point it at DuplicationLadderRecorder::finish instead. main ships this broken link because its ci-doc lacks --document-private-items; this branch checks private-item links, so the rustdoc gate fails without the fix.
The v0.7.0 sync bumped the five newly-added workspace crates (fgumi-fmt, fgumi-cli-common, fgumi-cli-macros, fgumi-pipeline-core, fgumi-pipeline-io) from 0.6.0 to the 0.7.0 workspace version in Cargo.toml during the rebase onto origin/main; refresh their entries in Cargo.lock to match.
|
Important Review skippedAuto reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
Note Reviews pausedUse the following commands to manage reviews:
Use the checkboxes below for quick actions:
Comment |
|
@coderabbitai pause |
✅ Action performedReviews paused. |
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #870 +/- ##
==========================================
- Coverage 94.68% 94.55% -0.13%
==========================================
Files 193 268 +75
Lines 119865 141664 +21799
==========================================
+ Hits 113490 133950 +20460
- Misses 6375 7714 +1339 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@coderabbitai review |
|
|
@coderabbitai review |
Apply the actionable review findings surfaced by the #870 review split (#877 foundation/arena-sort half, #878 pipeline-tree/step-ports half). ci / manifests: - bound the loom job with timeout-minutes, matching every sibling job - centralize the shared quote dependency in [workspace.dependencies]; syn stays per-crate (2.x vs 3.x in xtask, which one entry cannot represent) - constrain the workspace log pin to 0.4 and make fgumi-pipeline-io inherit it via workspace = true (every other crate already does); drop its redundant fgumi-bam-io / fgumi-bgzf / tempfile dev-dependencies fgumi-bgzf / fgumi-sort: - route the Vec decompress path through is_stored_block so the shared predicate cannot drift between entry points - reject zero workers in the arena benchmark; move the StagingBuffer constructor doc and #[must_use] back onto new; cover the malformed-frame NUL re-append in RawQuerynameKey::read_from; recycle drained BGZF block buffers back into the compressor; derive FRAME_DECOMP_CAP from worker_pool::ZSTD_FRAME_DECOMP_CAP so producer and consumer cannot drift - create spill files exclusively (create_new) in SyncSpillWriter, never truncating an existing chunk, matching SpillWrite::open_file fgumi-pipeline-core: - panic loudly in TypedStep2::build_fused_output_set rather than silently collapsing the load-bearing merge reorder to None; add the documented out-of-range assert in build_driver_storage (with a covering test); gate the fused stall test's should_panic on debug_assertions so the release test build reaches the is_err path; tighten both sticky-slot collision assertions to the full "loser gets no slot" contract - sync the deadlock_timeout_secs field doc with the now-optional stats contract (liveness comes from LivenessCounter) fgumi-pipeline-io: - name the where_only trybuild fixture's rejection arm (and regenerate its .stderr for the shifted span) - bound the input-derived BAM record block_size in both boundary scanners so a corrupt length prefix fails closed instead of buffering toward 4 GiB - honor PipelineReaderOpts::async_reader for the reopened regular-file block stream, matching read_bam_stdin - debug-assert the from_parsed range invariant at the constructor boundary and document it; fix the spill_gather framing test to treat frame_one_block's return as the next index, not a consumed count; add positive coverage for the drained-finish BGZF_EOF path in the WriteBgzfFile sink fgumi library: - fall back to the name-only group key for a mapped record with an empty CIGAR (both compute_group_key_from_raw copies) so the i32::MAX unclipped sentinel never merges distinct templates, matching the secondary path - guard cache_umi_position's aux slice with get() against a malformed record - carry legacy_overlap_window through the Codec -> CodecOptions projection - correct the step-library module doc: grouping steps mint their own serials
Apply the actionable review findings surfaced by the #870 review split (#877 foundation/arena-sort half, #878 pipeline-tree/step-ports half). ci / manifests: - bound the loom job with timeout-minutes, matching every sibling job - centralize the shared quote dependency in [workspace.dependencies]; syn stays per-crate (2.x vs 3.x in xtask, which one entry cannot represent) - constrain the workspace log pin to 0.4 and make fgumi-pipeline-io inherit it via workspace = true (every other crate already does); drop its redundant fgumi-bam-io / fgumi-bgzf / tempfile dev-dependencies fgumi-bgzf / fgumi-sort: - route the Vec decompress path through is_stored_block so the shared predicate cannot drift between entry points - reject zero workers in the arena benchmark; move the StagingBuffer constructor doc and #[must_use] back onto new; cover the malformed-frame NUL re-append in RawQuerynameKey::read_from; recycle drained BGZF block buffers back into the compressor; derive FRAME_DECOMP_CAP from worker_pool::ZSTD_FRAME_DECOMP_CAP so producer and consumer cannot drift - create spill files exclusively (create_new) in SyncSpillWriter, never truncating an existing chunk, matching SpillWrite::open_file fgumi-pipeline-core: - panic loudly in TypedStep2::build_fused_output_set rather than silently collapsing the load-bearing merge reorder to None; add the documented out-of-range assert in build_driver_storage (with a covering test); gate the fused stall test's should_panic on debug_assertions so the release test build reaches the is_err path; tighten both sticky-slot collision assertions to the full "loser gets no slot" contract - sync the deadlock_timeout_secs field doc with the now-optional stats contract (liveness comes from LivenessCounter) fgumi-pipeline-io: - name the where_only trybuild fixture's rejection arm (and regenerate its .stderr for the shifted span) - bound the input-derived BAM record block_size in both boundary scanners so a corrupt length prefix fails closed instead of buffering toward 4 GiB - honor PipelineReaderOpts::async_reader for the reopened regular-file block stream, matching read_bam_stdin - debug-assert the from_parsed range invariant at the constructor boundary and document it; fix the spill_gather framing test to treat frame_one_block's return as the next index, not a consumed count; add positive coverage for the drained-finish BGZF_EOF path in the WriteBgzfFile sink fgumi library: - fall back to the name-only group key for a mapped record with an empty CIGAR (both compute_group_key_from_raw copies) so the i32::MAX unclipped sentinel never merges distinct templates, matching the secondary path - guard cache_umi_position's aux slice with get() against a malformed record - carry legacy_overlap_window through the Codec -> CodecOptions projection - correct the step-library module doc: grouping steps mint their own serials
fgumi-bgzf: - bound the footer-claimed uncompressed size against MAX_UNCOMPRESSED_BLOCK_SIZE inside decompress_and_verify, so a corrupt ISIZE fails closed on the raw-Vec entry path (not only the slice path) before the resize zero-fills up to 4 GiB; cover it with a giant-ISIZE regression test on decompress_block_into_opts fgumi-sort: - stage each spill in a temp file and rename it into place only after finish() succeeds, so a crashed/errored write never leaves a partial file a streaming reader would accept as a shorter spill; persist_noclobber keeps the fail-closed (never-truncate) guarantee the create_new open gave. Cover atomic publish + no-clobber on both codecs - reuse the compressed-frame buffer in ZstdSpillWriter (compress_to_buffer) instead of letting compress() allocate a Vec per frame - correct the NAME_INLINE_CAP=44 rationale: SmallVec<[u8;40]> is a 48-byte NameBuf and [u8;44] a 56-byte one, so 44 costs 8 bytes/key over 40 (a deliberate trade to keep long Illumina names inline), not 'free' - anchor the sync-vs-pooled parity test to the original records first fgumi-pipeline-core: - cover both is_fusible_chain rejection guards (a Step2 merge with input_arity>1, and an output branch wired to an earlier-or-equal step) ci / fgumi: - pin the loom job's sccache-action to v0.0.11, matching every other job - cover cache_umi_position's out-of-range aux_offset guard with a malformed record whose inflated n_cigar_op pushes the derived aux offset past the body
…ed cancel The arena slot-fill decompressor (decompress_into_slice) short-circuited on uncompressed_size == 0, so a CRC-corrupted BGZF EOF marker (ISIZE 0 but no longer the exact BGZF_EOF bytes) returned Ok(0) without CRC verification. The Vec siblings skip only the exact BGZF_EOF marker and CRC-check every other zero-size block; gate the slot path the same way so a corrupted zero-size block is rejected rather than waved through. Adds a regression test with a correctly sized zero-length slot so the rejection comes from CRC, not slot sizing. Also add a fused single-thread driver test covering the documented PipelineError::Cancelled exit: a step cancels the shared signal on its first dispatch and the run must map to Cancelled with the sink short of the full item count, pinning the top-of-loop is_done() break.
The // Tests banner was left in place when the slice-fill decompression section was grafted onto the reader; it sat before ~290 lines of production code rather than the real #[cfg(test)] mod tests (line ~960), which is misleading. Remove it.
Foundation-half fixes from CodeRabbit's post-force-push full review: - fused.rs: drop steps before ChainContexts. run_fused_single_thread takes steps by value (a parameter) and builds contexts as a local; parameters drop after locals, which inverts the typed-handle cache invariant that every step is dropped before the contexts it cached from. Drop steps explicitly to restore it, keeping the perf cache. Documented in CLAUDE.md. - driver.rs / liveness.rs: correct the per-worker liveness sharding comments — a dedicated driver reuses worker_slot 0 and so shares slot 0 with pool worker 0; the atomic fetch_add still counts a coincident bump correctly. - tests.rs: strip the dotted suffix from Cargo dependency keys (log.workspace -> log) so the doc-enumeration test never reports a false undocumented dependency. - cli-macros: emit a blank doc line before the generated required-field notice so clap keeps it in long help rather than the short summary; add a help-visibility assertion for visible vs hidden aliases; reject memory multiplication overflow in the test parser via checked_mul. - miri.yml: anchor the unsafe_code location guards to attribute lines (a prose mention no longer fails the job) and send ::error:: annotations to stdout so GitHub renders them. - .coderabbit.yaml: split the pipeline path-instruction — the legacy src/lib homes keep the OrderedQueue/BamPipelineState guidance; the fgumi-pipeline-core crate gets its own entry naming its real symbols (queues.rs / ByteBoundedQueue), which the legacy text does not define. - sort-phase2-unification-deferral.md: the streaming slot's reorder buffer and in-flight counter are back for the block-parallel path; qualify the plain-FIFO and passive-slot claims to the inline path. The struct-level cfg_attr guard-set finding was a false positive: rustc expands a container cfg_attr before the multi_options macro runs, so a hidden clap attr is either caught as a literal or is inert on the inactive platform — unlike a field cfg_attr, which reaches the macro verbatim.
The round-5 hardening dropped `steps` explicitly before returning, but a panic unwinding out of `try_run_erased` or the stall `debug_assert!` skips that drop and reverts to natural order (locals before by-value parameters), which frees `ChainContexts` while the typed-handle cache in the steps still borrows from it. Re-bind `steps` as a local declared after `contexts` so reverse-declaration order runs steps-before-contexts on every exit, unwinding included, and drop the now-redundant explicit `drop(steps)`.
Two coverage gaps CodeRabbit flagged on the fused driver: - is_fusible_rejects_a_second_source: is_fusible_chain has four structural guards; the single-source rejection (a second source at index >= 1) was the only one without a test. The other three already have one. - drive_records_stats_for_success_and_error_dispatches: the driver's stats block records every dispatch when handed Some(stats), but every other fused test runs with None. One erroring chain exercises both arms — the source's Progress via record, then the mid's failure via record_error.
Lands the typed-step pipeline foundation onto
main, rebased on v0.7.0. Each unit below was developed and reviewed as its own sub-PR on a shared integration branch; this brings them tomaintogether so the workspace stays buildable at every crate boundary.The domain-crate split (
fgumi-dna,fgumi-umi,fgumi-consensus,fgumi-sort,fgumi-bgzf,fgumi-bam-io,fgumi-raw-bam, …) already shipped in v0.7.0 and is not introduced here — this PR adds the pipeline framework and a few supporting leaf crates on top of it.New crates (5):
fgumi-fmt,fgumi-cli-common,fgumi-cli-macros— format/CLI leaf crates (#672, #674);fgumi-pipeline-core,fgumi-pipeline-io— the typed-step pipeline framework (#697, #732).Typed-step pipeline foundation (
src/lib/pipeline/**+ the two pipeline crates): the typed-step pipeline tree and foundational steps (#735), read-side source steps (#736), grouping and closure-driven mid-steps (#743), the UMI-correction step (#821), per-stage CLI option projection (#744), and the shared grouping/library-lookup domain types (#734, #741). The chain-builder layer that consumes this framework is not part of this PR; the framework is dormant on every production command path here.Sort-engine additions to the existing
fgumi-sortcrate: the unified arena sort engine for all four sort orders (#720), the phase-2 merge handoff slot (#718), and BGZF slice decompression with caller-driven buffer recycling infgumi-bgzf(#714).Constituent PRs: #672, #674, #697, #714, #718, #720, #722, #732, #734, #735, #736, #741, #743, #744, #745, #821. Behavior is unchanged from what those merged; the only new work here is the v0.7.0 rebase and the crate-version reconciliation in
Cargo.lock.Regression validation (v0.7.0 → this branch, t=16)
Output is identical and there is no speed regression — measured directly, v0.7.0 binary vs this branch:
sort— byte-identical output (3,689,310 and 8,224,560 records).group(adjacency / edit / paired) — MI-equivalent (~86K molecules).dedup— content-identical (2,586,791 records).simplex— content-identical (81,768 consensus reads).sort/dedup.Not yet run — merge gate
The fgbio-correctness equivalency and WES/WGS-scale matrix (fgumi-benchmarks, AWS) has not been run. That is the release/promotion gate. Do not merge until it passes — the validation above establishes no regression versus v0.7.0, not cross-tool correctness at scale.