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};
#[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<()>>,
)> {
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() {
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))
}