use std::cell::RefCell;
use std::collections::HashMap;
use std::rc::Rc;
use std::time::Duration;
use gloo_timers::future::sleep;
use js_sys::{Function, Promise};
use mbus_client::services::ClientServices;
use mbus_client::services::coil::Coils;
#[cfg(feature = "file-record")]
use mbus_client::services::file_record::SubRequest;
#[cfg(feature = "file-record")]
use mbus_core::data_unit::common::MAX_PDU_DATA_LEN;
use mbus_core::errors::MbusError;
#[cfg(feature = "diagnostics")]
use mbus_core::function_codes::public::DiagnosticSubFunction;
#[cfg(feature = "diagnostics")]
use mbus_core::models::diagnostic::{ObjectId, ReadDeviceIdCode};
use mbus_core::transport::{
BackoffStrategy, JitterStrategy, ModbusConfig, ModbusTcpConfig, UnitIdOrSlaveAddr,
};
use wasm_bindgen::prelude::*;
use wasm_bindgen_futures::spawn_local;
use super::app::{PendingHandle, PendingMap, WasmAppRouter};
use mbus_network::WasmWsTransport;
const PIPELINE: usize = 10;
type Inner = ClientServices<WasmWsTransport, WasmAppRouter, PIPELINE>;
#[wasm_bindgen]
pub struct WasmModbusClient {
inner: Rc<RefCell<Inner>>,
pending: PendingMap,
unit_id: u8,
next_txn: u16,
}
#[wasm_bindgen]
impl WasmModbusClient {
#[wasm_bindgen(constructor)]
pub fn new(
ws_url: &str,
unit_id: u8,
response_timeout_ms: u32,
retry_attempts: u8,
tick_interval_ms: u32,
) -> Result<WasmModbusClient, JsValue> {
let pending: PendingMap = Rc::new(RefCell::new(HashMap::new()));
let app = WasmAppRouter::new(pending.clone());
let transport = WasmWsTransport::new(ws_url);
let config = ModbusConfig::Tcp(ModbusTcpConfig {
host: heapless::String::try_from("wasm")
.map_err(|_| JsValue::from_str("host string overflow"))?,
port: 0,
connection_timeout_ms: 5000,
response_timeout_ms,
retry_attempts,
retry_backoff_strategy: BackoffStrategy::Immediate,
retry_jitter_strategy: JitterStrategy::None,
retry_random_fn: None,
});
let mut inner_client = ClientServices::new(transport, app, config)
.map_err(|e| JsValue::from_str(&format!("{:?}", e)))?;
let _ = inner_client.connect();
let inner = Rc::new(RefCell::new(inner_client));
let weak = Rc::downgrade(&inner);
let tick_ms = tick_interval_ms as u64;
let idle_ms = core::cmp::max(50, tick_ms.saturating_mul(5));
spawn_local(async move {
loop {
match weak.upgrade() {
Some(rc) => {
let should_poll = {
let client = rc.borrow();
client.is_connected() && client.has_pending_requests()
};
if should_poll {
rc.borrow_mut().poll();
sleep(Duration::from_millis(tick_ms)).await;
} else {
sleep(Duration::from_millis(idle_ms)).await;
}
continue;
}
None => break, }
}
});
Ok(WasmModbusClient {
inner,
pending,
unit_id,
next_txn: 1,
})
}
pub fn is_connected(&self) -> bool {
self.inner.borrow().is_connected()
}
pub fn has_pending_requests(&self) -> bool {
self.inner.borrow().has_pending_requests()
}
pub fn reconnect(&mut self) -> bool {
for (_, handle) in self.pending.borrow_mut().drain() {
let _ = handle
.reject
.call1(&JsValue::NULL, &JsValue::from_str("ConnectionLost"));
}
self.inner.borrow_mut().reconnect().is_ok()
}
pub fn read_coils(&mut self, address: u16, quantity: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.coils()
.read_multiple_coils(txn_id, unit_addr, address, quantity);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn write_single_coil(&mut self, address: u16, value: bool) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.coils()
.write_single_coil(txn_id, unit_addr, address, value);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn read_holding_registers(&mut self, address: u16, quantity: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.read_holding_registers(txn_id, unit_addr, address, quantity);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn read_input_registers(&mut self, address: u16, quantity: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.read_input_registers(txn_id, unit_addr, address, quantity);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn write_single_register(&mut self, address: u16, value: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.write_single_register(txn_id, unit_addr, address, value);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn write_multiple_registers(
&mut self,
address: u16,
quantity: u16,
values: &[u16],
) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.write_multiple_registers(txn_id, unit_addr, address, quantity, values);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn read_discrete_inputs(&mut self, address: u16, quantity: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.discrete_inputs()
.read_discrete_inputs(txn_id, unit_addr, address, quantity);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
}
impl WasmModbusClient {
fn alloc_txn(&mut self) -> u16 {
let id = self.next_txn;
self.next_txn = self.next_txn.wrapping_add(1).max(1);
id
}
fn reject_immediate(&self, txn_id: u16, error: MbusError) {
if let Some(handle) = self.pending.borrow_mut().remove(&txn_id) {
let _ = handle
.reject
.call1(&JsValue::NULL, &JsValue::from_str(&format!("{:?}", error)));
}
}
}
fn make_promise() -> (Promise, Function, Function) {
let resolve_holder: Rc<RefCell<Option<Function>>> = Rc::new(RefCell::new(None));
let reject_holder: Rc<RefCell<Option<Function>>> = Rc::new(RefCell::new(None));
let r = resolve_holder.clone();
let rj = reject_holder.clone();
let promise = Promise::new(&mut move |res, rej| {
*r.borrow_mut() = Some(res);
*rj.borrow_mut() = Some(rej);
});
let resolve = resolve_holder.borrow_mut().take().unwrap();
let reject = reject_holder.borrow_mut().take().unwrap();
(promise, resolve, reject)
}
#[wasm_bindgen]
impl WasmModbusClient {
pub fn read_single_coil(&mut self, address: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.coils()
.read_single_coil(txn_id, unit_addr, address);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn write_multiple_coils(
&mut self,
address: u16,
quantity: u16,
packed_bytes: &[u8],
) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let coils_result =
Coils::new(address, quantity).and_then(|c| c.with_values(packed_bytes, quantity));
let result = match coils_result {
Ok(coils) => self
.inner
.borrow_mut()
.coils()
.write_multiple_coils(txn_id, unit_addr, address, &coils),
Err(e) => Err(e),
};
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn read_single_holding_register(&mut self, address: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.read_single_holding_register(txn_id, unit_addr, address);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn read_single_input_register(&mut self, address: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.read_single_input_register(txn_id, unit_addr, address);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn read_write_multiple_registers(
&mut self,
read_address: u16,
read_quantity: u16,
write_address: u16,
_write_quantity: u16,
values: &[u16],
) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.read_write_multiple_registers(
txn_id,
unit_addr,
read_address,
read_quantity,
write_address,
values,
);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn mask_write_register(&mut self, address: u16, and_mask: u16, or_mask: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.registers()
.mask_write_register(txn_id, unit_addr, address, and_mask, or_mask);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
pub fn read_single_discrete_input(&mut self, address: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.discrete_inputs()
.read_single_discrete_input(txn_id, unit_addr, address);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "fifo")]
pub fn read_fifo_queue(&mut self, address: u16) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.fifo()
.read_fifo_queue(txn_id, unit_addr, address);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "file-record")]
pub fn read_file_record(
&mut self,
file_number: u16,
record_number: u16,
record_length: u16,
) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let mut sub_req = SubRequest::new();
let result = sub_req
.add_read_sub_request(file_number, record_number, record_length)
.and_then(|_| {
self.inner
.borrow_mut()
.file_records()
.read_file_record(txn_id, unit_addr, &sub_req)
});
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "file-record")]
pub fn write_file_record(
&mut self,
file_number: u16,
record_number: u16,
values: &[u16],
) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let record_length = values.len() as u16;
let mut hv = heapless::Vec::<u16, MAX_PDU_DATA_LEN>::new();
for &v in values {
if hv.push(v).is_err() {
self.reject_immediate(txn_id, MbusError::BufferTooSmall);
return promise;
}
}
let mut sub_req = SubRequest::new();
let result = sub_req
.add_write_sub_request(file_number, record_number, record_length, hv)
.and_then(|_| {
self.inner
.borrow_mut()
.file_records()
.write_file_record(txn_id, unit_addr, &sub_req)
});
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "diagnostics")]
pub fn read_exception_status(&mut self) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.diagnostic()
.read_exception_status(txn_id, unit_addr);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "diagnostics")]
pub fn diagnostics(&mut self, sub_function: u16, data: &[u16]) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = DiagnosticSubFunction::try_from(sub_function)
.map_err(|_| MbusError::ReservedSubFunction(sub_function))
.and_then(|sf| {
self.inner
.borrow_mut()
.diagnostic()
.diagnostics(txn_id, unit_addr, sf, data)
});
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "diagnostics")]
pub fn get_comm_event_counter(&mut self) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.diagnostic()
.get_comm_event_counter(txn_id, unit_addr);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "diagnostics")]
pub fn get_comm_event_log(&mut self) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.diagnostic()
.get_comm_event_log(txn_id, unit_addr);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "diagnostics")]
pub fn report_server_id(&mut self) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = self
.inner
.borrow_mut()
.diagnostic()
.report_server_id(txn_id, unit_addr);
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
#[cfg(feature = "diagnostics")]
pub fn read_device_identification(
&mut self,
read_device_id_code: u8,
object_id: u8,
) -> Promise {
let txn_id = self.alloc_txn();
let (promise, resolve, reject) = make_promise();
self.pending
.borrow_mut()
.insert(txn_id, PendingHandle { resolve, reject });
let unit_addr = UnitIdOrSlaveAddr::new(self.unit_id).unwrap_or_default();
let result = ReadDeviceIdCode::try_from(read_device_id_code)
.map_err(|_| MbusError::InvalidDeviceIdCode)
.and_then(|code| {
self.inner
.borrow_mut()
.diagnostic()
.read_device_identification(txn_id, unit_addr, code, ObjectId::from(object_id))
});
if let Err(e) = result {
self.reject_immediate(txn_id, e);
}
promise
}
}