use rust_htslib::bam::Read;
use std::{format, thread};
use crate::binning::FinalResult;
use crate::config::NUM_BINS;
use crate::pbar::{DEFAULT_INTERVAL, get_spin_pb};
use crate::range::Range;
use crate::record::RecordData;
use crate::writer::write_final_result_csv;
#[allow(dead_code)]
pub fn run_pipeline(bam_path: &str, config: &super::PipelineConfig, ref_range: &Range<usize>, out_dir: &str, bam_name: &str) {
let worker_cnt = config.worker_count.max(1);
let bam_path = bam_path.to_string();
let length_min = 0;
let length_max = usize::MAX;
let (reader_sender, reader_recv) = crossbeam::channel::bounded(1000);
thread::scope(|thread_scope| {
thread_scope.spawn(move || {
let mut reader = rust_htslib::bam::Reader::from_path(&bam_path)
.expect(&format!("read {} error", bam_path));
for record in reader.records() {
if let Ok(record) = record {
if record.seq_len() >= length_min && record.seq_len() < length_max {
match reader_sender.send(record) {
Ok(_) => {}
Err(e) => {
tracing::error!("reader send error. {:#?}", e);
}
}
}
} else {
tracing::warn!("error record");
}
}
});
let (worker_sender, worker_recv) = crossbeam::channel::bounded(1000);
for _ in 0..worker_cnt {
thread_scope.spawn({
let reader_recv = reader_recv.clone();
let worker_sender = worker_sender.clone();
move || {
let mut result: FinalResult<NUM_BINS> = FinalResult::new_all();
for record in reader_recv {
if let Some(record_data) =
RecordData::from_record_within_ref_range(&record, ref_range)
{
result
.get_total_hist_mut()
.collect_ar_dw(record_data.ar.iter(), record_data.dw.iter());
result.get_total_per_base_hist_mut().collect_ar_dw(
record_data.bases.as_str(),
&record_data.ar,
&record_data.dw,
);
}
}
worker_sender.send(result)
}
});
}
drop(reader_recv);
drop(worker_sender);
let mut final_result: FinalResult<NUM_BINS> = FinalResult::new_all();
let pb = get_spin_pb(format!("processing"), DEFAULT_INTERVAL);
for result in worker_recv {
final_result.merge(&result);
pb.inc(1);
}
pb.finish();
write_final_result_csv(out_dir, bam_name, &final_result).unwrap();
});
}