use super::api::*;
use super::*;
use crate::{
daq::error::SiggenSnafu,
rt::PPM,
siggen::{self, Siggen, SiggenCommand, SiggenError, SourceDescriptor},
*,
};
use api::DaqApiMethods;
use api::*;
use array_init::from_iter;
use core::time;
use crossbeam::channel::{Receiver, Sender, TrySendError, bounded, unbounded};
use dasp_sample::Sample;
use snafu::prelude::*;
use std::{
any::Any,
collections::HashMap,
mem::{replace, swap},
sync::{
Arc, LazyLock, Mutex, Weak,
atomic::{AtomicBool, Ordering},
},
thread::sleep,
time::Duration,
};
use streamcmd::InputStreamCommand;
use streamdata::*;
use streammetadata::*;
use streammgr_details::*;
use streammsg::*;
use thread_priority::{ThreadPriority, set_current_thread_priority};
type Result<T> = std::result::Result<T, StreamMgrError>;
#[cfg(not(feature = "test_features"))]
static STREAMMGR_CREATED: AtomicBool = AtomicBool::new(false);
pub type SharedInQueue = Sender<InStreamMsg>;
pub type InQueues = Vec<SharedInQueue>;
pub type ApiMap = HashMap<DaqApiDescriptor, Box<dyn DaqApiMethods>>;
pub type DeviceList = Vec<DeviceInfo>;
#[derive(Debug)]
enum InputStreamState {
Undefined,
NotRunning {
queues: InQueues,
},
Running {
stream: Box<dyn Stream>,
stream_thread: JoinHandle<InQueues>,
commtx: Sender<InputStreamCommand>,
commrx: Receiver<std::result::Result<(), StreamMgrError>>,
},
RunningInDuplex {
stream: Box<dyn Stream>,
stream_thread: JoinHandle<(InQueues, Siggen, InQueues)>,
commtx: Sender<StreamCommand>,
commrx: Receiver<std::result::Result<(), StreamMgrError>>,
},
}
#[derive(Debug)]
#[allow(clippy::large_enum_variant)]
enum OutputStreamState {
Undefined,
NotRunning {
monqueues: InQueues,
siggen: Siggen,
},
RunningInDuplex,
Running {
stream: Box<dyn Stream>,
siggen_thread: JoinHandle<(Siggen, InQueues)>,
commtx: Sender<OutputStreamCommand>,
commrx: Receiver<std::result::Result<(), StreamMgrError>>,
},
}
#[cfg_attr(feature = "python-bindings", gen_stub_pyclass, pyclass(unsendable))]
#[derive(Debug)]
pub struct StreamMgr {
srcdesc: SourceDescriptor,
PPMmon: Weak<PPM>,
PPMinp: Weak<PPM>,
devs: DeviceList,
input_stream: InputStreamState,
output_stream: OutputStreamState,
apis: ApiMap,
#[allow(clippy::type_complexity)]
devices_scan: Option<(
Arc<AtomicBool>,
std::thread::JoinHandle<(ApiMap, DeviceList)>,
)>,
}
#[cfg(feature = "python-bindings")]
#[cfg_attr(feature = "python-bindings", gen_stub_pymethods, pymethods)]
impl StreamMgr {
#[new]
fn py_new() -> StreamMgr {
StreamMgr::new()
}
fn __repr__(&self) -> String {
format!("{self:#?}")
}
#[pyo3(name = "startDefaultInputStream")]
fn startDefaultInputStream_py(&mut self) -> PyResult<()> {
Ok(self.startDefaultInputStream()?)
}
#[pyo3(name = "startDefaultOutputStream")]
fn startDefaultOutputStream_py(&mut self) -> PyResult<()> {
Ok(self.startDefaultOutputStream()?)
}
#[pyo3(name = "startStream")]
fn startStream_py(&mut self, st: StreamType, d: &DaqConfig) -> PyResult<()> {
Ok(self.startStream(st, d)?)
}
#[pyo3(name = "stopStream")]
fn stopStream_py(&mut self, st: StreamType) -> PyResult<()> {
Ok(self.stopStream(st)?)
}
#[pyo3(name = "getDeviceInfo")]
fn getDeviceInfo_py(&mut self) -> PyResult<Vec<DeviceInfo>> {
Ok(self.getDeviceInfo())
}
#[pyo3(name = "getStatus")]
fn getStatus_py(&self, dir: StreamDirection) -> StreamStatus {
self.getStatus(dir)
}
#[pyo3(name = "getStreamMetaData")]
fn getStreamMetaData_py(&self, dir: StreamDirection) -> Option<StreamMetaData> {
self.getStreamMetaData(dir).map(|b| (*b).clone())
}
#[pyo3(name = "reScanDevices")]
fn reScanDevices_py(&mut self) -> PyResult<()> {
self.reScanDevices()?;
Ok(())
}
#[pyo3(name = "isSomeStreamRunning")]
fn isSomeStreamRunning_py(&self) -> bool {
self.isSomeStreamRunning()
}
#[pyo3(name = "isStreamRunning")]
fn isStreamRunning_py(&self, dir: StreamDirection) -> bool {
self.isStreamRunning(dir)
}
#[pyo3(name = "isStreamRunningOK")]
fn isStreamRunningOK_py(&self, dir: StreamDirection) -> PyResult<bool> {
Ok(self.isStreamRunningOK(dir))
}
#[pyo3(name = "siggenCommand")]
fn siggenCommand_py(&mut self, cmd: SiggenCommand) -> PyResult<()> {
Ok(self.siggenCommand(cmd)?)
}
}
impl Default for StreamMgr {
fn default() -> Self {
Self::new()
}
}
impl StreamMgr {
pub fn new() -> StreamMgr {
cfg_select! {
not(feature = "test_features")=> {
if STREAMMGR_CREATED
.compare_exchange(false, true, Ordering::Acquire, Ordering::Relaxed)
.is_err()
{
panic!("BUG: Stream manager is supposed to be a singleton");
}
},
_ => {}
}
let mut apis: HashMap<_, Box<dyn DaqApiMethods>> = HashMap::new();
#[cfg(feature = "cpal-api")]
apis.insert(DaqApiDescriptor::Cpal, Box::new(CpalApi::new()));
#[cfg(feature = "loopback-api")]
apis.insert(DaqApiDescriptor::Loopback, Box::new(LoopbackApi::new()));
#[cfg(feature = "uldaq-api")]
apis.insert(DaqApiDescriptor::Uldaq, Box::new(UldaqApi::new()));
let srcdesc = SourceDescriptor::Silence {};
let mut smgr = StreamMgr {
apis,
devs: vec![],
input_stream: InputStreamState::NotRunning { queues: vec![] },
output_stream: OutputStreamState::NotRunning {
monqueues: vec![],
siggen: Siggen::new(1, srcdesc.clone()),
},
PPMmon: Weak::new(),
PPMinp: Weak::new(),
devices_scan: None,
srcdesc,
};
smgr.reScanDevices().unwrap();
smgr
}
pub fn new_with_devices() -> Self {
let mut smgr = Self::new();
while smgr.isDeviceScanRunning() {
std::thread::sleep(Duration::from_millis(100));
smgr.updateDeviceListFromScan();
}
smgr
}
pub fn reScanDevices(&mut self) -> Result<()> {
ensure!(
self.devices_scan.is_none(),
DeviceScanAlreadyInProgressSnafu
);
ensure!(
matches!(self.input_stream, InputStreamState::NotRunning { .. }),
InputStreamAlreadyRunningSnafu
);
ensure!(
matches!(self.output_stream, OutputStreamState::NotRunning { .. }),
OutputStreamAlreadyRunningSnafu
);
let mut apis = HashMap::new();
swap(&mut self.apis, &mut apis);
let scanfinished = Arc::new(AtomicBool::new(false));
let scanfinished_clone = scanfinished.clone();
self.devices_scan = Some((
scanfinished_clone,
std::thread::spawn(move || {
create_thread_pool_if_not_created();
cfg_select! {
all(not(debug_assertions), target_os = "linux") => {
let print_gag = gag::Gag::stderr();
if let Err(e) = print_gag {
eprintln!("Unable to capture stderr: {e}");
}
},
target_os = "linux" => {
eprintln!("****** Any possible ALSA errors printed below are suppressed in release builds. ******");
},
_ => {}
}
let mut all_devices = vec![];
for api in apis.values() {
let devs = api.getDeviceInfo();
if let Ok(devs) = devs {
all_devices.extend(devs);
}
}
#[cfg(target_os = "linux")]
eprintln!("****** End of any possible ALSA errors. ******");
scanfinished.store(true, Ordering::Relaxed);
(apis, all_devices)
}),
));
Ok(())
}
pub fn isSomeStreamRunning(&self) -> bool {
!matches!(self.input_stream, InputStreamState::NotRunning { .. })
|| !matches!(self.output_stream, OutputStreamState::NotRunning { .. })
}
pub fn isStreamRunning(&self, dir: StreamDirection) -> bool {
match dir {
StreamDirection::Input => match &self.input_stream {
InputStreamState::NotRunning { .. } => false,
InputStreamState::Running { .. } => true,
InputStreamState::RunningInDuplex { .. } => true,
InputStreamState::Undefined => unreachable!(),
},
StreamDirection::Output => match &self.output_stream {
OutputStreamState::NotRunning { .. } => false,
OutputStreamState::Running { .. } => true,
OutputStreamState::RunningInDuplex => true,
OutputStreamState::Undefined => unreachable!(),
},
}
}
pub fn isStreamRunningOK(&self, dir: StreamDirection) -> bool {
match dir {
StreamDirection::Input => match &self.input_stream {
InputStreamState::NotRunning { .. } => false,
InputStreamState::Running { stream, .. } => {
matches!(stream.status(dir), StreamStatus::Running { .. })
}
InputStreamState::RunningInDuplex { stream, .. } => {
matches!(stream.status(dir), StreamStatus::Running { .. })
}
InputStreamState::Undefined => unreachable!(),
},
StreamDirection::Output => match &self.output_stream {
OutputStreamState::NotRunning { .. } => false,
OutputStreamState::Running { stream, .. } => {
matches!(stream.status(dir), StreamStatus::Running { .. })
}
OutputStreamState::RunningInDuplex => {
let InputStreamState::RunningInDuplex { stream, .. } = &self.input_stream
else {
unreachable!(
"Invalid input stream state, does not match output stream state, which is in duplex mode"
)
};
matches!(stream.status(dir), StreamStatus::Running { .. })
}
OutputStreamState::Undefined => unreachable!(),
},
}
}
pub fn getStreamMetaData(&self, dir: StreamDirection) -> Option<Arc<StreamMetaData>> {
match dir {
StreamDirection::Input => match &self.input_stream {
InputStreamState::NotRunning { .. } => None,
InputStreamState::Running { stream, .. } => stream.inMetaData(),
InputStreamState::RunningInDuplex { stream, .. } => stream.inMetaData(),
InputStreamState::Undefined => unreachable!(),
},
StreamDirection::Output => {
if let InputStreamState::RunningInDuplex { stream, .. } = &self.input_stream {
stream.outMetaData()
} else {
match &self.output_stream {
OutputStreamState::Undefined => unreachable!(),
OutputStreamState::NotRunning { .. } => None,
OutputStreamState::RunningInDuplex => unreachable!(),
OutputStreamState::Running { stream, .. } => stream.outMetaData(),
}
}
}
}
}
pub fn getStatus(&self, dir: StreamDirection) -> StreamStatus {
if let InputStreamState::RunningInDuplex { stream, .. } = &self.input_stream {
return stream.status(dir);
}
match dir {
StreamDirection::Input => {
match &self.input_stream {
InputStreamState::NotRunning { .. } => StreamStatus::NotRunning {},
InputStreamState::Running { stream, .. } => {
return stream.status(dir);
}
InputStreamState::RunningInDuplex { .. } => {
unreachable!()
}
InputStreamState::Undefined => unreachable!(),
};
}
StreamDirection::Output => {
match &self.output_stream {
OutputStreamState::NotRunning { .. } => StreamStatus::NotRunning {},
OutputStreamState::RunningInDuplex => unreachable!(),
OutputStreamState::Running { stream, .. } => {
return stream.status(dir);
}
OutputStreamState::Undefined => unreachable!(),
};
}
}
StreamStatus::NotRunning {}
}
pub fn getSourceDescriptor(&self) -> SourceDescriptor {
self.srcdesc.clone()
}
pub fn setSiggenSource(&mut self, src: SourceDescriptor) -> Result<()> {
if let InputStreamState::RunningInDuplex { commtx, commrx, .. } = &self.input_stream {
commtx
.send(StreamCommand::OutputStreamCommand(
OutputStreamCommand::SiggenCommand(SiggenCommand::ChangeSource {
src: src.clone(),
}),
))
.unwrap();
match commrx.recv().unwrap() {
Ok(()) => {
self.srcdesc = src;
return Ok(());
}
err @ Err(_) => {
return err;
}
}
}
match &mut self.output_stream {
OutputStreamState::NotRunning {
monqueues: _,
siggen,
} => {
siggen
.applyCommand(SiggenCommand::ChangeSource { src: src.clone() })
.context(SiggenSnafu {})?;
self.srcdesc = src;
Ok(())
}
OutputStreamState::RunningInDuplex => {
unreachable!("Cannot get here!")
}
OutputStreamState::Running { commtx, commrx, .. } => {
commtx
.send(OutputStreamCommand::SiggenCommand(
SiggenCommand::ChangeSource { src: src.clone() },
))
.unwrap();
match commrx.recv().unwrap() {
Ok(_) => {
self.srcdesc = src;
Ok(())
}
err @ Err(_) => err,
}
}
OutputStreamState::Undefined => unreachable!(),
}
}
fn updateDeviceListFromScan(&mut self) {
if let Some((scan_done, joinhandle)) = self.devices_scan.take() {
if scan_done.load(Ordering::Relaxed) {
let (apilist, devlist) = joinhandle.join().expect("Device scan panicked");
self.apis = apilist;
self.devs = devlist;
self.devices_scan = None;
} else {
self.devices_scan = Some((scan_done, joinhandle));
}
}
}
pub fn isDeviceScanRunning(&self) -> bool {
self.devices_scan.is_some()
}
pub fn getDeviceInfo(&mut self) -> Vec<DeviceInfo> {
self.updateDeviceListFromScan();
self.devs.clone()
}
pub fn addInQueue(&mut self, tx: Sender<InStreamMsg>) {
match &mut self.input_stream {
InputStreamState::NotRunning { queues } => {
queues.push(tx);
}
InputStreamState::Running { commtx, commrx, .. } => {
commtx.send(InputStreamCommand::AddInQueue(tx)).unwrap();
commrx
.recv()
.unwrap()
.expect("Adding a queue should never fail")
}
InputStreamState::RunningInDuplex { commtx, commrx, .. } => {
commtx
.send(StreamCommand::InputStreamCommand(
InputStreamCommand::AddInQueue(tx),
))
.unwrap();
commrx
.recv()
.unwrap()
.expect("Adding a queue should never fail")
}
InputStreamState::Undefined => unreachable!(),
}
}
pub fn addMonitorQueue(&mut self, tx: SharedInQueue) {
if let InputStreamState::RunningInDuplex { commtx, commrx, .. } = &self.input_stream {
commtx
.send(StreamCommand::OutputStreamCommand(
OutputStreamCommand::AddMonitorQueue(tx),
))
.unwrap();
commrx
.recv()
.unwrap()
.expect("Adding a queue should never fail");
return;
};
match &mut self.output_stream {
OutputStreamState::NotRunning { monqueues, .. } => {
monqueues.push(tx);
}
OutputStreamState::RunningInDuplex => unreachable!(),
OutputStreamState::Running {
stream: _,
siggen_thread: _,
commtx,
commrx,
} => {
commtx
.send(OutputStreamCommand::AddMonitorQueue(tx))
.unwrap();
commrx
.recv()
.unwrap()
.expect("Adding a queue should never fail");
}
OutputStreamState::Undefined => unreachable!(),
}
}
fn find_device(&self, cfg: &DaqConfig) -> Result<&DeviceInfo> {
ensure!(
self.devices_scan.is_none(),
DeviceScanAlreadyInProgressSnafu
);
if let Some(matching_dev) = self
.devs
.iter()
.find(|&d| d.device_name == cfg.device_name && d.api == cfg.api)
{
return Ok(matching_dev);
}
DeviceNotAvailableSnafu {
device_name: &cfg.device_name,
}
.fail()
}
pub fn startStream(&mut self, stype: StreamType, cfg: &DaqConfig) -> Result<()> {
self.updateDeviceListFromScan();
ensure!(
self.devices_scan.is_none(),
DeviceScanAlreadyInProgressSnafu
);
match stype {
StreamType::Input | StreamType::Duplex => {
self.startInputOrDuplexStream(stype, cfg)?;
}
StreamType::Output => {
self.startOutputStream(cfg)?;
}
}
Ok(())
}
fn startOutputStream(&mut self, cfg: &DaqConfig) -> Result<()> {
let stream = replace(&mut self.output_stream, OutputStreamState::Undefined);
match stream {
OutputStreamState::NotRunning { monqueues, siggen } => {
let (tx, rx): (Sender<Arc<RawStreamData>>, Receiver<Arc<RawStreamData>>) =
unbounded();
let startstream = |rx, tx, mut siggen: Siggen, mon_queues| -> Result<_> {
let api = self
.getDaqApi(&cfg.api)
.with_context(|| ApiNotAvailableSnafu {
apiname: cfg.api.name(),
})?;
let devinfo = self.find_device(cfg)?;
let stream = api.startOutputStream(devinfo, cfg, rx)?;
let meta = stream.outMetaData().expect("No stream metadata available");
siggen.reset(meta.samplerate).context(SiggenSnafu)?;
let (siggen_thread, commtx, commrx) =
startSiggenThread(meta, siggen, tx, mon_queues)?;
Ok((stream, siggen_thread, commtx, commrx))
};
match startstream(rx, tx, siggen.clone(), monqueues.clone()) {
Ok((stream, siggen_thread, commtx, commrx)) => {
self.output_stream = OutputStreamState::Running {
stream,
siggen_thread,
commtx,
commrx,
}
}
Err(e) => {
self.output_stream = OutputStreamState::NotRunning { monqueues, siggen };
return Err(e);
}
}
}
_ => {
let und = replace(&mut self.output_stream, stream);
assert!(matches!(und, OutputStreamState::Undefined));
return OutputStreamAlreadyRunningSnafu.fail();
}
}
Ok(())
}
fn startInputOrDuplexStream(&mut self, stype: StreamType, cfg: &DaqConfig) -> Result<()> {
assert!(!matches!(self.input_stream, InputStreamState::Undefined));
ensure!(
cfg.numberEnabledInChannels() > 0,
DAQConfigSnafu {
msg: "At least one input channel should be enabled \
for an input stream"
}
);
let duplex = matches!(stype, StreamType::Duplex);
ensure!(
!(stype == StreamType::Duplex && cfg.numberEnabledOutChannels() == 0),
DAQConfigSnafu {
msg: "At least one output channel should be enabled for a duplex stream"
}
);
ensure!(
!(matches!(self.output_stream, OutputStreamState::Running { .. }) && duplex),
DAQConfigSnafu {
msg: "An output stream is already running. Please first stop existing output stream."
}
);
let instream = replace(&mut self.input_stream, InputStreamState::Undefined);
match instream {
InputStreamState::NotRunning { queues } => {
if duplex {
let ostream = replace(&mut self.output_stream, OutputStreamState::Undefined);
let startduplexstream = |mut iqueues: InQueues,
siggen,
mut mon_queues: InQueues|
-> Result<_> {
assert!(!duplex);
let (tx, rx_in): (Sender<InStreamMsg>, Receiver<InStreamMsg>) = unbounded();
let (tx_out, rx_out) = unbounded();
let api = self.getDaqApi(&cfg.api).context(ApiNotAvailableSnafu {
apiname: cfg.api.name(),
})?;
let devinfo = self.find_device(cfg)?;
let stream =
api.startInputOrDuplexStream(stype, devinfo, cfg, tx, Some(rx_out))?;
let in_meta = stream
.inMetaData()
.expect("No input stream metadata available");
let out_meta = stream
.outMetaData()
.expect("No output stream metadata for duplex stream!");
sendMsgToAllQueuesRemoveUnused(
&mut iqueues,
InStreamMsg::StreamStarted(
in_meta.clone(),
PreCaptureBuffer::NotLoaded,
),
);
sendMsgToAllQueuesRemoveUnused(
&mut mon_queues,
InStreamMsg::StreamStarted(
out_meta.clone(),
PreCaptureBuffer::NotLoaded,
),
);
let (threadhandle, commtx, commrx) = startDuplexThread(
in_meta, rx_in, iqueues, out_meta, siggen, tx_out, mon_queues,
)?;
Ok((stream, threadhandle, commtx, commrx))
};
let (monqueues, siggen) =
if let OutputStreamState::NotRunning { monqueues, siggen } = ostream {
(monqueues, siggen)
} else {
unreachable!()
};
match startduplexstream(queues.clone(), siggen.clone(), monqueues.clone()) {
Ok((stream, stream_thread, commtx, commrx)) => {
let _ = replace(
&mut self.input_stream,
InputStreamState::RunningInDuplex {
stream,
stream_thread,
commtx,
commrx,
},
);
let _ = replace(
&mut self.output_stream,
OutputStreamState::RunningInDuplex {},
);
Ok(())
}
Err(e) => {
let _ = replace(
&mut self.input_stream,
InputStreamState::NotRunning { queues },
);
let _ = replace(
&mut self.output_stream,
OutputStreamState::NotRunning { monqueues, siggen },
);
Err(e)
}
}
} else {
let startinstream = |mut iqueues: InQueues| -> Result<_> {
assert!(!duplex);
let (tx, rx): (Sender<InStreamMsg>, Receiver<InStreamMsg>) = unbounded();
let api = self.getDaqApi(&cfg.api).context(ApiNotAvailableSnafu {
apiname: cfg.api.name(),
})?;
let devinfo = self.find_device(cfg)?;
let stream = api.startInputOrDuplexStream(stype, devinfo, cfg, tx, None)?;
let meta = stream
.inMetaData()
.expect("No input stream metadata available");
sendMsgToAllQueuesRemoveUnused(
&mut iqueues,
InStreamMsg::StreamStarted(meta.clone(), PreCaptureBuffer::NotLoaded),
);
let (threadhandle, commtx, commrx) =
startInputStreamThread(meta, rx, iqueues);
Ok((stream, threadhandle, commtx, commrx))
};
match startinstream(queues.clone()) {
Ok((stream, stream_thread, commtx, commrx)) => {
self.input_stream = InputStreamState::Running {
stream,
stream_thread,
commtx,
commrx,
};
Ok(())
}
Err(e) => {
let und = replace(
&mut self.input_stream,
InputStreamState::NotRunning { queues },
);
assert!(matches!(und, InputStreamState::Undefined));
Err(e)
}
}
}
}
_ => {
let _ = replace(&mut self.input_stream, instream);
InputStreamAlreadyRunningSnafu.fail()
}
}
}
pub fn startDefaultInputStream(&mut self) -> Result<()> {
self.updateDeviceListFromScan();
while self.isDeviceScanRunning() {
eprintln!("Cannot yet start stream: a device scan is still in progress.");
sleep(Duration::from_millis(20));
self.updateDeviceListFromScan();
}
cfg_select! {
feature = "cpal-api" => {
let stream = replace(&mut self.input_stream, InputStreamState::Undefined);
if let InputStreamState::NotRunning { queues } = stream {
let startstream = |mut iqueues: InQueues| -> Result<_> {
let (tx, rx): (Sender<InStreamMsg>, Receiver<InStreamMsg>) = unbounded();
let cpal_api: &CpalApi = self
.getDaqApiT::<CpalApi>()
.context(CPALNotAvailableSnafu)?;
let stream = cpal_api.startDefaultInputStream(tx)?;
let meta = stream
.inMetaData()
.expect("No input stream metadata available");
sendMsgToAllQueuesRemoveUnused(&mut iqueues, InStreamMsg::StreamStarted(meta.clone(), PreCaptureBuffer::NotLoaded));
let (threadhandle, commtx, commrx) = startInputStreamThread(meta, rx, iqueues);
Ok((stream, threadhandle, commtx, commrx))
};
match startstream(queues.clone()) {
Ok((stream, stream_thread, commtx, commrx)) => {
self.input_stream = InputStreamState::Running { stream, stream_thread, commtx, commrx };
Ok(())
}
Err(e) => {
let _ = replace(
&mut self.input_stream,
InputStreamState::NotRunning { queues },
);
Err(e)
}
}
} else {
let _ = replace(&mut self.input_stream, stream);
InputStreamAlreadyRunningSnafu.fail()
}
},
_ => {
CPALNotAvailableSnafu.fail()
}
}
}
pub fn startDefaultOutputStream(&mut self) -> Result<()> {
self.updateDeviceListFromScan();
while self.isDeviceScanRunning() {
eprintln!("Cannot yet start stream: a device scan is still in progress.");
sleep(Duration::from_millis(20));
self.updateDeviceListFromScan();
}
cfg_select! {
feature = "cpal-api" => {
let stream = replace(&mut self.output_stream, OutputStreamState::Undefined);
if let OutputStreamState::NotRunning { monqueues, siggen } = stream {
let startstream = |mut mon_queues: InQueues, siggen| -> Result<_> {
let (tx, rx) = unbounded();
let cpal_api: &CpalApi =
self.getDaqApiT::<CpalApi>().expect("CPal API not present");
let stream = cpal_api.startDefaultOutputStream(rx)?;
let meta = stream.outMetaData().expect("Output metadata not available");
sendMsgToAllQueuesRemoveUnused(&mut mon_queues, InStreamMsg::StreamStarted(meta.clone(), PreCaptureBuffer::NotLoaded));
Ok((stream, startSiggenThread(meta, siggen, tx, mon_queues)?))
};
match startstream(monqueues.clone(), siggen.clone()) {
Ok((stream, (siggen_thread, commtx, commrx))) => {
let _ = replace(&mut self.output_stream, OutputStreamState::Running { stream, siggen_thread, commtx, commrx });
Ok(())
},
Err(e) => {
let _ = replace(
&mut self.output_stream,
OutputStreamState::NotRunning { monqueues, siggen }
);
Err(e)
},
}
} else {
let _ = replace(&mut self.output_stream, stream);
OutputStreamAlreadyRunningSnafu.fail()
}
}, _ => {
CPALNotAvailableSnafu.fail()
}
} }
pub fn stopInputStream(&mut self) -> Result<()> {
assert!(!matches!(self.input_stream, InputStreamState::Undefined));
let stream = replace(&mut self.input_stream, InputStreamState::Undefined);
if let InputStreamState::Running {
stream: _,
stream_thread,
commtx,
commrx,
} = stream
{
commtx.send(InputStreamCommand::StopThread).unwrap();
let _ = commrx.recv().unwrap();
let queues = stream_thread.join();
self.input_stream = InputStreamState::NotRunning { queues };
Ok(())
} else if let InputStreamState::RunningInDuplex {
stream: _,
stream_thread,
commtx,
commrx,
} = stream
{
commtx.send(StreamCommand::StopThread).unwrap();
let _ = commrx.recv().unwrap();
let (iqueues, siggen, monqueues) = stream_thread.join();
self.input_stream = InputStreamState::NotRunning { queues: iqueues };
self.output_stream = OutputStreamState::NotRunning { monqueues, siggen };
Ok(())
} else {
self.input_stream = stream;
InputStreamNotRunningSnafu.fail()
}
}
pub fn stopOutputStream(&mut self) -> Result<()> {
let stream = replace(&mut self.output_stream, OutputStreamState::Undefined);
if let OutputStreamState::Running {
stream: _,
siggen_thread,
commtx,
commrx,
} = stream
{
commtx.send(OutputStreamCommand::StopThread).unwrap();
let _ = commrx.recv().unwrap();
let (siggen, monqueues) = siggen_thread.join();
self.output_stream = OutputStreamState::NotRunning { monqueues, siggen };
Ok(())
} else {
let _ = replace(&mut self.output_stream, stream);
OutputStreamNotRunningSnafu.fail()
}
}
pub fn stopStream(&mut self, st: StreamType) -> Result<()> {
assert!(!matches!(self.input_stream, InputStreamState::Undefined));
assert!(!matches!(self.output_stream, OutputStreamState::Undefined));
match st {
StreamType::Input | StreamType::Duplex => self.stopInputStream(),
StreamType::Output => self.stopOutputStream(),
}
}
pub fn siggenCommand(&mut self, cmd: SiggenCommand) -> Result<()> {
if let InputStreamState::RunningInDuplex { commtx, commrx, .. } = &self.input_stream {
commtx
.send(StreamCommand::OutputStreamCommand(
OutputStreamCommand::SiggenCommand(cmd),
))
.unwrap();
commrx.recv().unwrap()
} else if let OutputStreamState::Running { commtx, commrx, .. } = &self.output_stream {
commtx
.send(OutputStreamCommand::SiggenCommand(cmd))
.unwrap();
commrx.recv().unwrap()
} else if let OutputStreamState::NotRunning {
monqueues: _,
siggen,
} = &mut self.output_stream
{
siggen.applyCommand(cmd).context(SiggenSnafu)
} else {
unreachable!()
}
}
pub fn getPPMInput(&self) -> Option<Arc<PPM>> {
self.PPMinp.upgrade()
}
pub fn setPPMInput(&mut self, ppmMon: &Arc<PPM>) {
if self.PPMinp.upgrade().is_some() {
panic!("Input PPM is already running!")
}
self.PPMinp = Arc::downgrade(ppmMon);
}
pub fn getDaqApi(&self, apidescr: &DaqApiDescriptor) -> Option<&dyn DaqApiMethods> {
self.apis.get(apidescr).map(|v| &**v)
}
pub fn getDaqApiT<T>(&self) -> Option<&T>
where
T: DaqApiMethods,
{
for val in self.apis.values() {
let val: &dyn Any = val.as_any();
let val = val.downcast_ref();
if val.is_some() {
return val;
}
}
None
}
pub fn getPPMMon(&self) -> Option<Arc<PPM>> {
self.PPMmon.upgrade()
}
pub fn setPPMMon(&mut self, ppmMon: &Arc<PPM>) {
if self.PPMmon.upgrade().is_some() {
panic!("Monitor PPM is already running!")
}
self.PPMmon = Arc::downgrade(ppmMon);
}
} impl Drop for StreamMgr {
fn drop(&mut self) {
if matches!(
self.input_stream,
InputStreamState::Running { .. } | InputStreamState::RunningInDuplex { .. }
) {
let _ = self.stopInputStream();
}
if matches!(self.output_stream, OutputStreamState::Running { .. }) {
let _ = self.stopOutputStream();
}
while self.devices_scan.is_some() {
sleep(Duration::from_millis(1));
self.updateDeviceListFromScan();
}
cfg_select! {
not(feature = "test_features") => {
STREAMMGR_CREATED.store(false, Ordering::Release);
},
_ => {}
}
}
}