use std::fs::File;
use std::sync::Arc;
use std::path::PathBuf;
use bgzip::{write::BGZFMultiThreadWriter, BGZFError, Compression};
use crate::{command::getraw, fileformat::{shard::{CellID, ReadPair}, CellUMI}};
use crate::fileformat::ReadPairWriter;
use super::{bam, detect_fileformat::get_fq_filename_r2_from_r1, ConstructFromPath, StreamingReadPairReader};
use seq_io::fastq::Reader as FastqReader;
use seq_io::fastq::Record as FastqRecord;
type ListReadWithBarcode = Arc<(CellID,Arc<Vec<ReadPair>>)>;
#[derive(Debug,Clone)]
pub struct BascetPairedFastqWriterFactory {
}
impl BascetPairedFastqWriterFactory {
pub fn new() -> BascetPairedFastqWriterFactory {
BascetPairedFastqWriterFactory {}
}
}
impl ConstructFromPath<BascetPairedFastqWriter> for BascetPairedFastqWriterFactory {
fn new_from_path(&self, fname: &PathBuf) -> anyhow::Result<BascetPairedFastqWriter> { BascetPairedFastqWriter::new(fname)
}
}
pub struct BascetPairedFastqWriter {
pub writer_r1: BGZFMultiThreadWriter<File>,
pub writer_r2: BGZFMultiThreadWriter<File>
}
impl BascetPairedFastqWriter {
fn new(path: &PathBuf) -> anyhow::Result<BascetPairedFastqWriter>{
println!("starting writer for paired FASTQ {:?}", path);
let spath = path.to_string_lossy();
let last_pos = spath.rfind("R1");
if let Some(last_pos) = last_pos {
let mut spath_r2 = spath.as_bytes().to_vec();
spath_r2[last_pos+1] = b'2';
let path_r2 = String::from_utf8(spath_r2).unwrap();
let out_buffer_r1 = File::create(&path).expect("Failed to create fastq.gz output file");
let writer_r1 = BGZFMultiThreadWriter::new(out_buffer_r1, Compression::default());
let out_buffer_r2 = File::create(&path_r2).expect("Failed to create fastq.gz output file");
let writer_r2 = BGZFMultiThreadWriter::new(out_buffer_r2, Compression::default());
Ok(BascetPairedFastqWriter {
writer_r1: writer_r1,
writer_r2: writer_r2
})
} else {
anyhow::bail!("Could not find R2 for fastq file {:?}", path);
}
}
}
impl ReadPairWriter for BascetPairedFastqWriter {
fn write_reads_for_cell(&mut self, cell_id:&CellID, list_reads: &Arc<Vec<ReadPair>>) {
let mut read_num = 0;
for rp in list_reads.iter() {
write_paired_fastq_read(
&mut self.writer_r1,
&make_fastq_readname(read_num, &cell_id, &rp.umi, 1),
&rp.r1,
&rp.q1
).unwrap();
write_paired_fastq_read(
&mut self.writer_r2,
&make_fastq_readname(read_num, &cell_id, &rp.umi, 2),
&rp.r2,
&rp.q2
).unwrap();
read_num+=1;
}
}
fn writing_done(&mut self) -> anyhow::Result<()> {
anyhow::Ok(())
}
}
fn write_paired_fastq_read<W: std::io::Write>(
writer: &mut W,
head: &Vec<u8>,
seq:&Vec<u8>,
qual:&Vec<u8>
) -> Result<(), BGZFError> {
writer.write_all(b"@")?;
writer.write_all(head.as_slice())?;
writer.write_all(b"\n")?;
writer.write_all(seq.as_slice())?;
writer.write_all(b"\n+\n")?;
writer.write_all(&qual.as_slice())?;
writer.write_all(b"\n")?;
Ok(())
}
fn make_fastq_readname(
read_num: u32,
cell_id: &CellID,
cell_umi: &CellUMI,
illumna_read_index: u32
) -> Vec<u8> {
let name=format!("BASCET_{}:{}:{} {}",
cell_id,
String::from_utf8(cell_umi.clone()).unwrap(),
read_num,
illumna_read_index);
name.as_bytes().to_vec() }
pub struct PairedFastqStreamingReadPairReader {
forward_file: FastqReader<Box<dyn std::io::Read>>,
reverse_file: FastqReader<Box<dyn std::io::Read>>,
last_rp: Option<(Vec<u8>,ReadPair)>
}
impl PairedFastqStreamingReadPairReader {
pub fn new(fname: &PathBuf) -> anyhow::Result<PairedFastqStreamingReadPairReader> {
let fname_r2 = get_fq_filename_r2_from_r1(&fname).unwrap();
let mut forward_file = getraw::open_fastq(&fname).unwrap(); let mut reverse_file = getraw::open_fastq(&fname_r2).unwrap();
let r1 = forward_file.next();
let r2 = reverse_file.next();
let rp = if let Some(r1) = r1 {
let r1= r1.as_ref().expect("Error reading record r1");
let r2 = r2.unwrap().expect("Error reading record r2");
let (cell_id, umi) = bam::readname_to_cell_umi(r1.head());
Some((cell_id.to_vec(), ReadPair {
r1: r1.seq().to_vec(),
r2: r2.seq().to_vec(),
q1: r1.qual().to_vec(),
q2: r2.qual().to_vec(),
umi: umi.to_vec()
}))
} else {
println!("Warning: empty input BAM");
None
};
Ok(PairedFastqStreamingReadPairReader {
forward_file: forward_file,
reverse_file: reverse_file,
last_rp: rp
})
}
}
impl StreamingReadPairReader for PairedFastqStreamingReadPairReader {
fn get_reads_for_next_cell(
&mut self
) -> anyhow::Result<Option<ListReadWithBarcode>> {
if let Some((current_cell, last_rp)) = self.last_rp.clone() {
let mut reads:Vec<ReadPair> = Vec::new();
reads.push(last_rp);
self.last_rp = None;
while let Some(r1) = self.forward_file.next() {
let r1 = r1.expect("Error reading record r1");
let r2 = self.reverse_file.next().unwrap().expect("Error reading record r2");
let (cell_id, umi) = bam::readname_to_cell_umi(r1.head());
let rp = ReadPair {
r1: r1.seq().to_vec(),
r2: r2.seq().to_vec(),
q1: r1.qual().to_vec(),
q2: r2.qual().to_vec(),
umi: umi.to_vec()
};
if cell_id == current_cell {
reads.push(rp);
} else {
self.last_rp = Some((
cell_id.to_vec(),
rp
));
break;
}
}
let reads = Arc::new(reads);
let cellid_reads = (
String::from_utf8(current_cell).unwrap(),
reads
);
Ok(Some(Arc::new(cellid_reads)))
} else {
Ok(None)
}
}
}
#[derive(Debug,Clone)]
pub struct PairedFastqStreamingReadPairReaderFactory {
}
impl PairedFastqStreamingReadPairReaderFactory {
pub fn new() -> PairedFastqStreamingReadPairReaderFactory {
PairedFastqStreamingReadPairReaderFactory {}
}
}
impl ConstructFromPath<PairedFastqStreamingReadPairReader> for PairedFastqStreamingReadPairReaderFactory {
fn new_from_path(&self, fname: &PathBuf) -> anyhow::Result<PairedFastqStreamingReadPairReader> {
PairedFastqStreamingReadPairReader::new(fname)
}
}