md_analysis 0.2.1

molecular dynamics
/// 三阶段 worker 流水线:reader → stats workers → writer
///
/// 使用 std::sync primitives(Mutex + Condvar)实现 fan-out channel,
/// 解决标准库 mpsc 不支持多消费者(multiple receivers)的问题。
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)]
/// 运行三阶段并行流水线,返回 (Histogram, PerBaseHistograms)。
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();
    });
}