use super::super::*;
use crate::{
daq::error::StreamMgrError,
siggen::{self, Siggen, SiggenCommand, SiggenError},
tools::find_unused_buf,
*,
};
use crossbeam::channel::{Receiver, Sender, TrySendError, unbounded};
use dasp_sample::{Sample, ToSample};
use snafu::ResultExt;
use std::{
any::Any,
collections::{HashMap, VecDeque},
mem::replace,
sync::{Arc, Mutex, Weak, atomic::AtomicBool},
time::Duration,
};
type Result<T> = std::result::Result<T, StreamMgrError>;
#[allow(clippy::type_complexity)]
pub(crate) fn startSiggenThread(
meta: Arc<StreamMetaData>,
mut siggen: Siggen,
tx: Sender<Arc<RawStreamData>>,
mut mon_queues: Vec<SharedInQueue>,
) -> Result<(
JoinHandle<(Siggen, InQueues)>,
Sender<OutputStreamCommand>,
Receiver<Result<()>>,
)> {
let (commtx_res, commrx) = unbounded();
let (commtx, commrx_res) = unbounded();
let nchannels = meta.nchannels();
siggen.setAllMute(true);
if siggen.nchannels() != nchannels {
siggen.setNChannels(nchannels);
}
siggen.reset(meta.samplerate).context(SiggenSnafu)?;
let threadhandle = spawn(
move || {
let sleep_time_us = Duration::from_micros(
(0.1 * 1e6 * meta.framesPerBlock as Flt / *meta.samplerate) as u64,
);
let mut bufs: VecDeque<Arc<RawStreamData>> = VecDeque::with_capacity(10);
let mut floatbuf: Vec<Flt> = vec![0.; nchannels * meta.framesPerBlock];
let mut ctr = 0;
'infy: loop {
if let Ok(streamcommand) = commrx.try_recv() {
match streamcommand {
OutputStreamCommand::StopThread => {
commtx.send(Ok(())).unwrap();
mon_queues.retain(|q| q.send(InStreamMsg::StreamStopped).is_ok());
break 'infy;
}
OutputStreamCommand::SiggenCommand(cmd) => {
let res = siggen.applyCommand(cmd);
commtx
.send(res.map_err(|source| StreamMgrError::SiggenError { source }))
.unwrap();
}
OutputStreamCommand::AddMonitorQueue(tx) => {
if let Ok(()) = tx.send(InStreamMsg::StreamStarted(
meta.clone(),
PreCaptureBuffer::NotLoaded,
)) {
mon_queues.push(tx);
}
}
}
}
if tx.is_empty() {
siggen.genSignal(&mut floatbuf);
let mut buftouse = find_unused_buf(&mut bufs).
unwrap_or_else(||
match meta.rawDatatype {
DataType::I8 => Arc::new(RawStreamData::Datai8(vec![0; floatbuf.len()])),
DataType::F32 => Arc::new(RawStreamData::Dataf32(vec![0.; floatbuf.len()])),
DataType::F64 => Arc::new(RawStreamData::Dataf64(vec![0.; floatbuf.len()])),
DataType::I16 => Arc::new(RawStreamData::Datai16(vec![0; floatbuf.len()])),
DataType::I32 => Arc::new(RawStreamData::Datai32(vec![0; floatbuf.len()])),
DataType::I24 => Arc::new(RawStreamData::Datai24(vec![dasp_sample::I24::EQUILIBRIUM; floatbuf.len()])),
});
let mutbuf = Arc::get_mut(&mut buftouse)
.expect("Buffer taken tat is in use. Not possible");
match meta.rawDatatype {
DataType::I8 => {
if let RawStreamData::Datai8(v) = mutbuf {
v.iter_mut()
.zip(floatbuf.iter())
.for_each(|(v, f)| *v = (*f).to_sample_());
} else {
unreachable!("Buffer is of wrong type");
}
}
DataType::I16 => {
if let RawStreamData::Datai16(v) = mutbuf {
v.iter_mut()
.zip(floatbuf.iter())
.for_each(|(v, f)| *v = (*f).to_sample_());
} else {
unreachable!("Buffer is of wrong type");
}
}
DataType::I24 => {
if let RawStreamData::Datai32(v) = mutbuf {
v.iter_mut()
.zip(floatbuf.iter())
.for_each(|(v, f)| *v = (*f).to_sample_());
} else {
unreachable!("Buffer is of wrong type");
}
}
DataType::I32 => {
if let RawStreamData::Datai32(v) = mutbuf {
v.iter_mut()
.zip(floatbuf.iter())
.for_each(|(v, f)| *v = (*f).to_sample_());
} else {
unreachable!("Buffer is of wrong type");
}
}
DataType::F32 => {
if let RawStreamData::Dataf32(v) = mutbuf {
v.iter_mut()
.zip(floatbuf.iter())
.for_each(|(v, f)| *v = (*f).to_sample_());
} else {
unreachable!("Buffer is of wrong type");
}
}
DataType::F64 => {
if let RawStreamData::Dataf64(v) = mutbuf {
v.iter_mut()
.zip(floatbuf.iter())
.for_each(|(v, f)| *v = (*f).to_sample_());
} else {
unreachable!("Buffer is of wrong type");
}
}
}
bufs.push_back(buftouse.clone());
if let Err(_e) = tx.send(buftouse) {
}
if !mon_queues.is_empty() {
let floatdat = Array2::from_shape_fn(
(meta.framesPerBlock, meta.nchannels()).f(),
|(frame, channel)| floatbuf[frame * meta.nchannels() + channel],
);
let msg = InStreamMsg::InStreamData(Arc::new(
InStreamData::newFromConverted(ctr, meta.clone(), floatdat, false),
));
mon_queues.retain(|q| q.send(msg.clone()).is_ok());
}
ctr += 1;
} else {
}
}
std::thread::sleep(sleep_time_us);
(siggen, mon_queues)
},
ThreadPriority::Normal,
);
Ok((threadhandle, commtx_res, commrx_res))
}