use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex, mpsc};
use std::time::Duration;
use super::super::responder::WorkerStatus;
use super::super::scpi::ScpiClient;
use super::config::{CHANNEL_COUNT, ChannelConfig, Config};
use super::peripheral::{ChannelState, InstrumentState, SiglentSdg2042X};
const SAFE_LOAD: &str = "100000";
const STARTUP_QUERY_COUNT: u32 = 12;
const STARTUP_COMMAND_COUNT: u32 = 7;
const STARTUP_PROCESSING_MARGIN: Duration = Duration::from_millis(250);
const ESR_ERROR_MASK: u8 = 0b0011_1100;
struct State {
safe_state_pending: bool,
next: Option<InstrumentState>,
applied: InstrumentState,
status: WorkerStatus,
}
struct Inner {
config: Config,
state: Mutex<State>,
changed: Condvar,
}
pub struct SiglentSdg2042XDriver {
inner: Arc<Inner>,
}
impl SiglentSdg2042XDriver {
pub fn new(config: Config) -> Result<Self, String> {
config.validate()?;
Ok(Self {
inner: Arc::new(Inner {
config,
state: Mutex::new(State {
safe_state_pending: false,
next: None,
applied: InstrumentState::default(),
status: WorkerStatus::default(),
}),
changed: Condvar::new(),
}),
})
}
pub fn peripheral(&self) -> SiglentSdg2042X {
SiglentSdg2042X::new(self.inner.config.connection.serial_number)
}
pub fn identity(&self) -> Option<String> {
self.inner.state.lock().ok()?.status.identity()
}
pub(super) fn shared_handle(&self) -> Self {
Self {
inner: self.inner.clone(),
}
}
pub(super) fn startup_timeout(&self) -> Duration {
self.inner.config.connection.startup_timeout(
STARTUP_QUERY_COUNT,
STARTUP_COMMAND_COUNT,
STARTUP_PROCESSING_MARGIN,
)
}
pub(super) fn run_worker(
self,
stop: Arc<AtomicBool>,
startup: mpsc::SyncSender<Result<(), String>>,
) -> Result<(), String> {
siglent_worker(self, stop, startup)
}
pub(super) fn submit(&self, request: InstrumentState) {
let request = request.normalized(&self.inner.config.channels);
let mut state = self.inner.state.lock().unwrap();
state.next = Some(request);
self.inner.changed.notify_one();
}
pub(super) fn request_safe_state(&self) {
let mut state = self.inner.state.lock().unwrap();
state.safe_state_pending = true;
state.next = None;
self.inner.changed.notify_one();
}
pub(super) fn applied(&self) -> Result<InstrumentState, String> {
let state = self
.inner
.state
.lock()
.map_err(|_| "SDG2042X state poisoned")?;
if let Some(error) = state.status.error() {
Err(error)
} else {
Ok(state.applied)
}
}
pub(super) fn latched_error(&self) -> Option<String> {
self.inner.state.lock().ok()?.status.error()
}
#[cfg(test)]
pub(super) fn take_queued(&self) -> Option<InstrumentState> {
take_pending(&mut self.inner.state.lock().unwrap())
}
#[cfg(test)]
pub(super) fn queued(&self) -> Option<InstrumentState> {
self.inner.state.lock().unwrap().next
}
}
fn take_pending(state: &mut State) -> Option<InstrumentState> {
if state.safe_state_pending {
state.safe_state_pending = false;
Some(InstrumentState::default())
} else {
state.next.take()
}
}
fn siglent_worker(
driver: SiglentSdg2042XDriver,
stop: Arc<AtomicBool>,
startup: mpsc::SyncSender<Result<(), String>>,
) -> Result<(), String> {
let result = siglent_worker_inner(&driver, &stop, &startup);
if let Err(error) = &result
&& let Ok(mut state) = driver.inner.state.lock()
{
state.status.latch_error(format!("SDG2042X: {error}"));
}
result
}
fn siglent_worker_inner(
driver: &SiglentSdg2042XDriver,
stop: &Arc<AtomicBool>,
startup: &mpsc::SyncSender<Result<(), String>>,
) -> Result<(), String> {
let config = &driver.inner.config;
let mut client = match ScpiClient::connect(&config.connection) {
Ok(client) => client,
Err(err) => {
let _ = startup.send(Err(format!("SDG2042X connection failed: {err}")));
return Err(err);
}
};
let setup = setup_siglent(&mut client, config);
let identity = match setup {
Ok(identity) => identity,
Err(err) => {
let error = match client.shutdown() {
Ok(()) => err,
Err(shutdown_err) => {
format!("{err}; additionally failed to close SCPI transport: {shutdown_err}")
}
};
let _ = startup.send(Err(format!("SDG2042X setup failed: {error}")));
return Err(error);
}
};
driver
.inner
.state
.lock()
.unwrap()
.status
.set_identity(identity);
let _ = startup.send(Ok(()));
let run_result = loop {
let mut state = driver.inner.state.lock().unwrap();
while !state.safe_state_pending && state.next.is_none() && !stop.load(Ordering::Relaxed) {
state = driver
.inner
.changed
.wait_timeout(state, Duration::from_millis(20))
.unwrap()
.0;
}
if stop.load(Ordering::Relaxed) {
break Ok(());
}
let request = take_pending(&mut state).unwrap();
let applied = state.applied;
drop(state);
if let Err(err) = apply_request(&mut client, config, applied, request) {
break Err(format!("failed to apply commanded state: {err}"));
}
driver.inner.state.lock().unwrap().applied = request;
};
let shutdown_result = shutdown_siglent(&mut client);
combine_worker_results(run_result, shutdown_result)
}
fn shutdown_siglent(client: &mut ScpiClient) -> Result<(), String> {
let safe_result = safe_outputs(client);
let transport_result = client.shutdown();
match (safe_result, transport_result) {
(Err(err), Err(transport_err)) => Err(format!(
"safe-state failure: {err}; transport shutdown also failed: {transport_err}"
)),
(Err(err), Ok(())) => Err(format!("safe-state failure: {err}")),
(Ok(()), Err(err)) => Err(err),
(Ok(()), Ok(())) => Ok(()),
}
}
pub(super) fn combine_worker_results(
run_result: Result<(), String>,
shutdown_result: Result<(), String>,
) -> Result<(), String> {
match (run_result, shutdown_result) {
(Err(err), Err(shutdown_err)) => Err(format!(
"{err}; additionally failed to shut down SDG2042X: {shutdown_err}"
)),
(Err(err), Ok(())) => Err(err),
(Ok(()), Err(err)) => Err(format!("failed to shut down SDG2042X: {err}")),
(Ok(()), Ok(())) => Ok(()),
}
}
fn setup_siglent(client: &mut ScpiClient, config: &Config) -> Result<String, String> {
let identity = client.identify()?;
config.connection.validate_identity(&identity)?;
client.command("*CLS")?;
safe_outputs(client)?;
for (index, channel) in config.channels.iter().enumerate() {
let number = index + 1;
let load = channel.load.scpi();
client.command(&format!("C{number}:OUTP LOAD,{load}"))?;
verify_output_state(client, number, false, &load)?;
verify_safe_waveform(client, number)?;
}
verify_standard_event_status(client)?;
Ok(identity)
}
fn verify_standard_event_status(client: &mut ScpiClient) -> Result<(), String> {
let response = client.query("*ESR?")?;
let value = response
.split_ascii_whitespace()
.next_back()
.ok_or_else(|| "empty SDG2042X *ESR? response".to_owned())?
.parse::<u8>()
.map_err(|err| format!("invalid SDG2042X *ESR? response `{response}`: {err}"))?;
if value & ESR_ERROR_MASK == 0 {
Ok(())
} else {
Err(format!(
"SDG2042X setup set SCPI error bits in event status `{response}`"
))
}
}
fn safe_outputs(client: &mut ScpiClient) -> Result<(), String> {
let mut errors = Vec::new();
for number in 1..=CHANNEL_COUNT {
if let Err(err) = client.command(&safe_waveform_command(number)) {
errors.push(format!("channel {number} zero command: {err}"));
}
}
if let Err(err) = expect_operation_complete(client) {
errors.push(format!("zero completion: {err}"));
}
for number in 1..=CHANNEL_COUNT {
if let Err(err) = verify_safe_waveform(client, number) {
errors.push(format!("channel {number} zero readback: {err}"));
}
}
for number in 1..=CHANNEL_COUNT {
if let Err(err) = client.command(&format!("C{number}:OUTP OFF,LOAD,{SAFE_LOAD}")) {
errors.push(format!("channel {number} output-off command: {err}"));
}
}
if let Err(err) = expect_operation_complete(client) {
errors.push(format!("output-off completion: {err}"));
}
for number in 1..=CHANNEL_COUNT {
if let Err(err) = verify_output_state(client, number, false, SAFE_LOAD) {
errors.push(format!("channel {number} output-off readback: {err}"));
}
}
if errors.is_empty() {
Ok(())
} else {
Err(errors.join("; "))
}
}
fn command_safe_channel(client: &mut ScpiClient, number: usize) -> Result<(), String> {
client.command(&safe_waveform_command(number))?;
client.command(&format!("C{number}:OUTP OFF,LOAD,{SAFE_LOAD}"))
}
fn safe_waveform_command(channel_number: usize) -> String {
format!("C{channel_number}:BSWV WVTP,DC,OFST,0")
}
fn verify_safe_waveform(client: &mut ScpiClient, channel_number: usize) -> Result<(), String> {
let readback = client.query(&format!("C{channel_number}:BSWV?"))?;
let waveform = parameter_value(&readback, "WVTP");
let offset = parameter_value(&readback, "OFST").and_then(|value| {
value
.trim_end_matches(|c: char| c.is_ascii_alphabetic())
.parse()
.ok()
});
if waveform.is_some_and(|value| value.eq_ignore_ascii_case("DC")) && offset == Some(0.0) {
Ok(())
} else {
Err(format!(
"channel {channel_number} waveform readback `{readback}` was not 0 V DC"
))
}
}
fn verify_output_state(
client: &mut ScpiClient,
channel_number: usize,
enabled: bool,
load: &str,
) -> Result<(), String> {
let readback = client.query(&format!("C{channel_number}:OUTP?"))?;
let state = if enabled { "ON" } else { "OFF" };
let expected = format!("C{channel_number}:OUTP {state},LOAD,{load}");
if readback.to_ascii_uppercase().starts_with(&expected) {
Ok(())
} else {
Err(format!(
"channel {channel_number} output readback `{readback}` did not start with `{expected}`"
))
}
}
fn parameter_value<'a>(response: &'a str, name: &str) -> Option<&'a str> {
let mut tokens = response.split(',').map(str::trim);
while let Some(token) = tokens.next() {
if token
.split_ascii_whitespace()
.next_back()
.is_some_and(|token| token.eq_ignore_ascii_case(name))
{
return tokens.next();
}
}
None
}
pub(super) fn apply_request(
client: &mut ScpiClient,
config: &Config,
applied: InstrumentState,
desired: InstrumentState,
) -> Result<(), String> {
let mut changed = false;
for (index, (applied, desired)) in applied
.channels()
.into_iter()
.zip(desired.channels())
.enumerate()
{
if applied == desired {
continue;
}
let number = index + 1;
if desired.enabled == 0.0 {
if applied.enabled != 0.0 {
command_safe_channel(client, number)?;
changed = true;
}
continue;
}
let channel = &config.channels[index];
if applied.enabled == 0.0 {
client.command(&basic_wave_command(number, channel, desired))?;
let load = channel.load.scpi();
client.command(&format!("C{number}:OUTP ON,LOAD,{load}"))?;
changed = true;
} else if waveform_settings_changed(channel, applied, desired) {
client.command(&basic_wave_command(number, channel, desired))?;
changed = true;
}
}
if changed {
expect_operation_complete(client)?;
}
Ok(())
}
fn waveform_settings_changed(
config: &ChannelConfig,
applied: ChannelState,
desired: ChannelState,
) -> bool {
applied.offset_voltage_v != desired.offset_voltage_v
|| (config.waveform.uses_frequency() && applied.frequency_hz != desired.frequency_hz)
|| (config.waveform.uses_duty() && applied.pulse_duty_cycle != desired.pulse_duty_cycle)
|| (config.waveform.uses_phase() && applied.phase_deg != desired.phase_deg)
|| (!config.waveform.uses_offset() && applied.stdev != desired.stdev)
}
pub(super) fn basic_wave_command(
channel_number: usize,
config: &ChannelConfig,
request: ChannelState,
) -> String {
let waveform = config.waveform;
let mut command = format!("C{channel_number}:BSWV WVTP,{}", waveform.scpi());
if waveform.uses_frequency() {
command.push_str(&format!(",FRQ,{}", scpi_number(request.frequency_hz)));
}
if waveform.uses_amplitude() {
command.push_str(&format!(",AMP,{}", scpi_number(config.amplitude_vpp)));
}
if waveform.uses_offset() {
command.push_str(&format!(",OFST,{}", scpi_number(request.offset_voltage_v)));
} else {
command.push_str(&format!(",MEAN,{}", scpi_number(request.offset_voltage_v)));
command.push_str(&format!(",STDEV,{}", scpi_number(request.stdev)));
}
if waveform.uses_duty() {
command.push_str(&format!(
",DUTY,{}",
scpi_number(request.pulse_duty_cycle * 100.0)
));
}
if waveform.uses_phase() {
command.push_str(&format!(",PHSE,{}", scpi_number(request.phase_deg)));
}
command
}
fn expect_operation_complete(client: &mut ScpiClient) -> Result<(), String> {
let response = client.query("*OPC?")?;
if response.trim() == "1" {
Ok(())
} else {
Err(format!("unexpected *OPC? response `{response}`"))
}
}
pub(super) fn scpi_number(value: f64) -> String {
format!("{value:.17e}")
}