use std::path::PathBuf;
use anyhow::Context;
use futures::TryStreamExt;
use noodles::bam;
use noodles::bgzf;
use noodles::bgzf::writer::CompressionLevel;
use noodles::fasta;
use noodles::sam;
use tokio::fs::File;
use crate::utils::args::CompressionStrategy;
use crate::utils::args::NumberOfRecords;
use crate::utils::display::RecordCounter;
use crate::utils::formats;
use crate::utils::formats::cram::ParsedAsyncCRAMFile;
use crate::utils::formats::utils::IndexCheck;
pub async fn to_sam_async(
from: PathBuf,
to: PathBuf,
fasta: PathBuf,
max_records: NumberOfRecords,
) -> anyhow::Result<()> {
let ParsedAsyncCRAMFile {
mut reader, header, ..
} = formats::cram::open_and_parse_async(from, IndexCheck::None)
.await
.with_context(|| {
"opening and parsing CRAM file. Check that your reference FASTA matches \
the FASTA used to generate the file."
})?;
let repository = fasta::indexed_reader::Builder::default()
.build_from_path(fasta)
.map(fasta::repository::adapters::IndexedReader::new)
.map(fasta::Repository::new)
.with_context(|| "FASTA repository and associated index")?;
let handle = File::create(to).await?;
let mut writer = sam::AsyncWriter::new(handle);
writer.write_header(&header.parsed).await?;
let mut counter = RecordCounter::new();
let mut records = reader.records(&repository, &header.parsed);
while let Some(record) = records.try_next().await? {
let record = record.try_into_alignment_record(&header.parsed)?;
writer
.write_alignment_record(&header.parsed, &record)
.await?;
counter.inc();
if counter.time_to_break(&max_records) {
break;
}
}
Ok(())
}
pub async fn to_bam_async(
from: PathBuf,
to: PathBuf,
fasta: PathBuf,
max_records: NumberOfRecords,
compression_strategy: CompressionStrategy,
) -> anyhow::Result<()> {
let ParsedAsyncCRAMFile {
mut reader, header, ..
} = formats::cram::open_and_parse_async(from, IndexCheck::None)
.await
.with_context(|| {
"opening and parsing CRAM file. Check that your reference FASTA matches \
the FASTA used to generate the file."
})?;
let repository = fasta::indexed_reader::Builder::default()
.build_from_path(fasta)
.map(fasta::repository::adapters::IndexedReader::new)
.map(fasta::Repository::new)
.with_context(|| "FASTA repository and associated index")?;
let compression_level: CompressionLevel = compression_strategy.into();
let mut writer = File::create(to)
.await
.map(|f| {
bgzf::r#async::writer::Builder::default()
.set_compression_level(compression_level)
.build_with_writer(f)
})
.map(bam::AsyncWriter::from)
.with_context(|| "opening output filestream")?;
writer.write_header(&header.parsed).await?;
writer
.write_reference_sequences(header.parsed.reference_sequences())
.await?;
let mut counter = RecordCounter::new();
let mut records = reader.records(&repository, &header.parsed);
while let Some(record) = records.try_next().await? {
let record = record.try_into_alignment_record(&header.parsed)?;
writer
.write_alignment_record(&header.parsed, &record)
.await?;
counter.inc();
if counter.time_to_break(&max_records) {
break;
}
}
writer.shutdown().await?;
Ok(())
}