Skip to content
2 changes: 1 addition & 1 deletion crates/fgumi-bam-io/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,12 +27,12 @@ bstr = { workspace = true }
bytes = { workspace = true }
crossbeam-channel = { workspace = true }
log = { workspace = true }
tempfile = { workspace = true }

[target.'cfg(target_os = "linux")'.dependencies]
nix = { workspace = true }

[dev-dependencies]
tempfile = { workspace = true }
rstest = { workspace = true }
proptest = { workspace = true }
criterion = { workspace = true }
Expand Down
8 changes: 4 additions & 4 deletions crates/fgumi-bam-io/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,8 @@ pub use reader::{
};
pub use reorder::{DrainReady, ReorderBuffer};
pub use writer::{
BaiBuilder, BamWriter, BgzfWriterEnum, IndexingBamWriter, RawBamWriter, bai_sidecar_path,
create_bam_writer, create_indexing_bam_writer, create_optional_bam_writer,
create_raw_bam_writer, open_output_writer, write_bai_index, write_bai_sidecar,
write_bam_header,
AlignmentContext, BaiBuilder, BamWriter, BgzfWriterEnum, IndexingBamWriter, RawBamWriter,
bai_sidecar_path, create_bam_writer, create_indexing_bam_writer, create_optional_bam_writer,
create_raw_bam_writer, extract_alignment_context, open_output_writer, write_bai_index,
write_bai_sidecar, write_bam_header,
};
287 changes: 254 additions & 33 deletions crates/fgumi-bam-io/src/writer.rs

Large diffs are not rendered by default.

379 changes: 372 additions & 7 deletions crates/fgumi-pipeline-io/src/sink/write_bgzf.rs

Large diffs are not rendered by default.

7 changes: 6 additions & 1 deletion crates/fgumi-pipeline-io/src/sink/write_raw.rs
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,12 @@ mod tests {
fn raw_bytes_block_exposes_payload() {
let plain = DecompressedBlock { batch_serial: 0, bytes: b"ACGT".to_vec() };
assert_eq!(RawBytesBlock::bytes(&plain), b"ACGT");
let bgzf = BgzfBlock { batch_serial: 0, bytes: b"\x1f\x8b".to_vec(), uncompressed_size: 4 };
let bgzf = BgzfBlock {
batch_serial: 0,
bytes: b"\x1f\x8b".to_vec(),
uncompressed_size: 4,
index: None,
};
assert_eq!(RawBytesBlock::bytes(&bgzf), b"\x1f\x8b");
}

Expand Down
7 changes: 7 additions & 0 deletions crates/fgumi-pipeline-io/src/sort/arena_ingest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1789,6 +1789,7 @@ mod tests {
batch_serial: 0,
bytes: blk0.clone(),
uncompressed_size: u32::try_from(isz0).unwrap(),
index: None,
})
.is_ok()
);
Expand All @@ -1797,6 +1798,7 @@ mod tests {
batch_serial: 1,
bytes: blk1.clone(),
uncompressed_size: u32::try_from(isz1).unwrap(),
index: None,
})
.is_ok()
);
Expand Down Expand Up @@ -1869,6 +1871,7 @@ mod tests {
batch_serial: 0,
bytes: blk0.clone(),
uncompressed_size: u32::try_from(isz0).unwrap(),
index: None,
})
.is_ok()
);
Expand All @@ -1877,6 +1880,7 @@ mod tests {
batch_serial: 1,
bytes: blk1.clone(),
uncompressed_size: u32::try_from(isz1).unwrap(),
index: None,
})
.is_ok()
);
Expand Down Expand Up @@ -1922,6 +1926,7 @@ mod tests {
batch_serial: 99,
bytes: probe_blk,
uncompressed_size: u32::try_from(probe_isz).unwrap(),
index: None,
});
assert!(result.is_err(), "pool must be exhausted while run 0 Arc is still in flight");
// A failed admit (ensure_arena returned false) must not leak state: the
Expand All @@ -1940,6 +1945,7 @@ mod tests {
batch_serial: 2,
bytes: blk2.clone(),
uncompressed_size: u32::try_from(isz2).unwrap(),
index: None,
})
.is_ok(),
"run 1 admit must succeed after pool release"
Expand All @@ -1949,6 +1955,7 @@ mod tests {
batch_serial: 3,
bytes: blk3.clone(),
uncompressed_size: u32::try_from(isz3).unwrap(),
index: None,
})
.is_ok(),
"run 1 second admit must succeed"
Expand Down
1 change: 1 addition & 0 deletions crates/fgumi-pipeline-io/src/sort/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1826,6 +1826,7 @@ fn bgzf_blocks_for(
batch_serial: i as u64,
bytes: blocks.remove(0).data,
uncompressed_size: u32::try_from(payload.len()).expect("payload fits u32"),
index: None,
}
})
.collect()
Expand Down
1 change: 1 addition & 0 deletions crates/fgumi-pipeline-io/src/source/read_bam.rs
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,7 @@ impl Step for ReadBgzfBlocks {
)
})?,
bytes: raw.data,
index: None,
});
}

Expand Down
99 changes: 96 additions & 3 deletions crates/fgumi-pipeline-io/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,57 @@
use fgumi_pipeline_core::{HeapSize, Ordered};
use fgumi_raw_bam::RawRecord;

// ─────────────────────────────────────────────────────────────────────────────
// BamIndexManifest / RecordIndexEntry — inline BAI sidecar metadata.
// ─────────────────────────────────────────────────────────────────────────────

/// One record's virtual-offset bookkeeping within a compressed [`BgzfBlock`],
/// as needed to build a BAI index inline with compression.
#[derive(Debug, Clone, Copy)]
pub struct RecordIndexEntry {
/// Offset of the record's 4-byte `block_size` length prefix, measured
/// from the start of the *uncompressed* bytes of the batch that produced
/// this block (i.e. the BGZF virtual offset's uncompressed-offset half,
/// before the physical block offset is added in).
pub uoffset: u32,
/// Total record length in bytes: `4` (the length prefix) plus the body
/// length.
pub len: u32,
/// The record's alignment context (reference id, start/end, mapped
/// flag), or `None` for a record with no reference/position (fully
/// unplaced unmapped reads are not indexable by coordinate).
pub ctx: Option<fgumi_bam_io::AlignmentContext>,
}

/// Sidecar metadata riding alongside a compressed [`BgzfBlock`], carrying
/// everything an inline BAI indexer needs to attribute records to virtual
/// offsets without re-parsing the compressed bytes.
///
/// # Invariants
///
/// - `phys_comp_len.iter().map(|&n| n as u64).sum::<u64>() == bytes.len() as
/// u64`, where `bytes` is the sibling [`BgzfBlock::bytes`] this manifest
/// describes: the physical lengths of the constituent BGZF blocks sum to
/// the total compressed byte count.
/// - `phys_comp_len.len() == (uncompressed_size as
/// usize).div_ceil(fgumi_bgzf::BGZF_MAX_BLOCK_SIZE)`, where
/// `uncompressed_size` is the sibling [`BgzfBlock::uncompressed_size`]:
/// one physical block per `BGZF_MAX_BLOCK_SIZE`-sized (or smaller, for the
/// final one) chunk of uncompressed input.
/// - `records` is in batch order, which is final coordinate order.
/// - Each entry's `uoffset` is the offset of the record's 4-byte length
/// prefix within the batch's uncompressed bytes; `len` is `4 +` the body
/// length.
#[derive(Debug, Clone)]
pub struct BamIndexManifest {
/// Physical (compressed) length of each constituent BGZF block, in the
/// order those blocks appear in the sibling [`BgzfBlock::bytes`].
pub phys_comp_len: Vec<u32>,
/// Per-record virtual-offset bookkeeping, in batch (= final coordinate)
/// order.
pub records: Vec<RecordIndexEntry>,
}

// ─────────────────────────────────────────────────────────────────────────────
// BgzfBlock — raw compressed BGZF block + read-order serial.
// ─────────────────────────────────────────────────────────────────────────────
Expand All @@ -23,13 +74,22 @@ pub struct BgzfBlock {
pub bytes: Vec<u8>,
/// Decompressed size, parsed from the BGZF block header.
pub uncompressed_size: u32,
/// Sidecar index metadata for this block's records, populated only on the
/// inline-BAI-indexed compress path (`None` for every other producer,
/// including sentinel/EOF blocks).
pub index: Option<Box<BamIndexManifest>>,
}

impl HeapSize for BgzfBlock {
fn heap_size(&self) -> usize {
// Byte-bounded queues budget on resident heap, so account for the full
// allocation (`capacity`), not just the populated prefix (`len`).
self.bytes.capacity()
let manifest = self.index.as_ref().map_or(0, |m| {
std::mem::size_of::<BamIndexManifest>()
+ m.phys_comp_len.capacity() * std::mem::size_of::<u32>()
+ m.records.capacity() * std::mem::size_of::<RecordIndexEntry>()
});
self.bytes.capacity() + manifest
}
}

Expand Down Expand Up @@ -228,11 +288,44 @@ mod tests {

#[test]
fn bgzf_block_heap_size_matches_bytes_capacity() {
let b = BgzfBlock { batch_serial: 0, bytes: vec![0u8; 1024], uncompressed_size: 4096 };
let b = BgzfBlock {
batch_serial: 0,
bytes: vec![0u8; 1024],
uncompressed_size: 4096,
index: None,
};
assert_eq!(b.heap_size(), 1024);
assert_eq!(b.ordinal(), 0);
}

#[test]
fn bgzf_block_index_defaults_none_and_heap_counts_manifest() {
let plain =
BgzfBlock { batch_serial: 0, bytes: vec![0u8; 100], uncompressed_size: 0, index: None };
assert!(plain.index.is_none());
assert_eq!(plain.heap_size(), plain.bytes.capacity());

let manifest = Box::new(BamIndexManifest {
phys_comp_len: vec![10, 20],
records: vec![RecordIndexEntry { uoffset: 0, len: 50, ctx: None }],
});
let indexed = BgzfBlock {
batch_serial: 1,
bytes: vec![0u8; 30],
uncompressed_size: 0,
index: Some(manifest),
};
let m = indexed.index.as_ref().expect("manifest present");
assert_eq!(
indexed.heap_size(),
indexed.bytes.capacity()
+ std::mem::size_of::<BamIndexManifest>()
+ m.phys_comp_len.capacity() * std::mem::size_of::<u32>()
+ m.records.capacity() * std::mem::size_of::<RecordIndexEntry>(),
"heap_size must account for the manifest allocation exactly"
);
}

#[test]
fn decompressed_block_heap_size_matches_bytes_capacity() {
let b = DecompressedBlock { batch_serial: 7, bytes: vec![0u8; 4096] };
Expand All @@ -250,7 +343,7 @@ mod tests {
assert!(bytes.capacity() >= 8192 && bytes.len() == 100);
let cap = bytes.capacity();

let block = BgzfBlock { batch_serial: 0, bytes, uncompressed_size: 0 };
let block = BgzfBlock { batch_serial: 0, bytes, uncompressed_size: 0, index: None };
assert_eq!(block.heap_size(), cap);

let mut backing = Vec::with_capacity(4096);
Expand Down
Loading
Loading