lasprs 0.14.1

Library for Acoustic Signal Processing (Rust edition, with optional Python bindings via pyo3)
use crate::daq::SiggenSnafu;
use crate::*;
use crate::{
    daq::{
        InQueues, InStreamMsg, InputStreamCommand, OutputStreamCommand, RawStreamData,
        SharedInQueue, StreamCommand, StreamMetaData, error::StreamMgrError,
    },
    siggen::{self, Siggen},
};
use core::time;
use crossbeam::channel::{Receiver, Sender, unbounded};
use snafu::prelude::*;
use std::sync::Arc;
type Result<T> = std::result::Result<T, StreamMgrError>;

use super::{siggenthread, startInputStreamThread, startSiggenThread};

/// Starts a duplex thread that handles both input and output streams at the same time.
///
/// Args:
///     in_meta: Arc<StreamMetaData>
///     rx: Receiver<InStreamMsg>
///     iqueues: InQueues
///     out_meta: Arc<StreamMetaData>
///     siggen: Siggen
///     tx: Sender<Arc<RawStreamData>>
///     mon_queues: Vec<SharedInQueue>
///
/// Returns:
///     (
///         JoinHandle<(InQueues, Siggen, InQueues)>,
///         Sender<StreamCommand>,
///         Receiver<Result<()>>,
///     )
#[allow(clippy::type_complexity)]
pub fn startDuplexThread(
    in_meta: Arc<StreamMetaData>,
    rx: Receiver<InStreamMsg>,
    iqueues: InQueues,
    out_meta: Arc<StreamMetaData>,
    siggen: Siggen,
    tx: Sender<Arc<RawStreamData>>,
    mon_queues: Vec<SharedInQueue>,
) -> Result<(
    JoinHandle<(InQueues, Siggen, InQueues)>,
    Sender<StreamCommand>,
    Receiver<Result<()>>,
)> {
    // Bi-directional communication between input stream thread and stream manager
    let (commtx_ret, commrx) = unbounded();
    let (commtx, commrx_ret) = unbounded();

    let (siggenthread, siggensender, siggenrx) =
        startSiggenThread(out_meta, siggen, tx, mon_queues)?;

    let threadhandle = spawn(
        move || {
            let (istream, istreamsender, istreamres) = startInputStreamThread(in_meta, rx, iqueues);
            'infy: loop {
                if let Ok(streammsg) = commrx.recv() {
                    // message obtained.
                    match streammsg {
                        StreamCommand::InputStreamCommand(input_stream_command) => {
                            istreamsender.send(input_stream_command).unwrap();
                            commtx.send(istreamres.recv().unwrap()).unwrap();
                        }
                        StreamCommand::OutputStreamCommand(output_stream_command) => {
                            siggensender.send(output_stream_command).unwrap();
                            commtx.send(siggenrx.recv().unwrap()).unwrap();
                        }
                        StreamCommand::StopThread => {
                            siggensender.send(OutputStreamCommand::StopThread).unwrap();
                            istreamsender.send(InputStreamCommand::StopThread).unwrap();
                            break 'infy;
                        }
                    }
                }
            }

            let iqueues = istream.join();
            let (siggen, mon_queues) = siggenthread.join();
            (iqueues, siggen, mon_queues)
        },
        ThreadPriority::High,
    );

    Ok((threadhandle, commtx_ret, commrx_ret))
}