use core::{
fmt::Debug,
sync::atomic::{AtomicU8, Ordering},
};
use alloc::vec::Vec;
use maybe_async::maybe_async;
use crate::{
application_protocol::{
application_pdu::ApplicationPdu,
confirmed::{
ComplexAck, ComplexAckService, ConfirmedRequest, ConfirmedRequestService, SimpleAck,
},
services::{
change_of_value::{CovNotification, SubscribeCov},
i_am::IAm,
read_property::{ReadProperty, ReadPropertyAck},
read_property_multiple::{ReadPropertyMultiple, ReadPropertyMultipleAck},
read_range::{ReadRange, ReadRangeAck},
time_synchronization::TimeSynchronization,
who_is::WhoIs,
write_property::WriteProperty,
write_property_multiple::WritePropertyMultiple,
},
unconfirmed::UnconfirmedRequest,
},
common::{
error::Error,
io::{Reader, Writer},
},
network_protocol::{
data_link::{DataLink, DataLinkFunction},
network_pdu::{DestinationAddress, MessagePriority, NetworkAddress, NetworkMessage, NetworkPdu},
},
};
#[derive(Debug)]
pub struct Bacnet<T>
where
T: NetworkIo + Debug,
{
pub io: T,
invoke_id: AtomicU8,
}
#[allow(async_fn_in_trait)]
#[cfg(feature = "defmt")]
#[maybe_async(AFIT)] pub trait NetworkIo {
type Error: Debug + defmt::Format;
async fn read(&self, buf: &mut [u8]) -> Result<usize, Self::Error>;
async fn write(&self, buf: &[u8]) -> Result<usize, Self::Error>;
}
#[cfg(not(feature = "defmt"))]
#[allow(async_fn_in_trait)]
#[maybe_async(AFIT)] pub trait NetworkIo {
type Error: Debug;
async fn read(&self, buf: &mut [u8]) -> Result<usize, Self::Error>;
async fn write(&self, buf: &[u8]) -> Result<usize, Self::Error>;
async fn disconnect(&self) -> Result<bool, Self::Error>;
}
#[derive(Debug)]
#[cfg_attr(feature = "defmt", derive(defmt::Format))]
pub enum BacnetError<T>
where
T: NetworkIo,
{
Io(T::Error),
Codec(Error),
InvokeId(InvokeIdError),
}
impl<T: NetworkIo> From<Error> for BacnetError<T> {
fn from(value: Error) -> Self {
Self::Codec(value)
}
}
#[derive(Debug)]
#[cfg_attr(feature = "defmt", derive(defmt::Format))]
pub struct InvokeIdError {
pub expected: u8,
pub actual: u8,
}
impl<T> Bacnet<T>
where
T: NetworkIo + Debug,
{
pub fn new(io: T) -> Self {
Self {
io,
invoke_id: AtomicU8::new(0),
}
}
pub fn into_inner(self) -> T {
self.io
}
#[maybe_async()]
pub async fn who_is(&self, buf: &mut [u8]) -> Result<Option<Vec<IAm>>, BacnetError<T>> {
let apdu = ApplicationPdu::UnconfirmedRequest(UnconfirmedRequest::WhoIs(WhoIs {}));
let dst = Some(DestinationAddress::new(0xffff, None));
let message = NetworkMessage::Apdu(apdu);
let npdu = NetworkPdu::new(None, dst, false, MessagePriority::Normal, message);
let data_link = DataLink::new(DataLinkFunction::OriginalBroadcastNpdu, Some(npdu));
let mut writer = Writer::new(buf);
data_link.encode(&mut writer);
let buffer = writer.to_bytes();
self.io.write(buffer).await.map_err(BacnetError::Io)?;
let mut iams:Vec<IAm> = Vec::new();
for _ in 0..5 {
match self.io.read(buf).await{
Ok(n) => {
let buf = &buf[..n];
let mut reader = Reader::default();
let message = DataLink::decode(&mut reader, buf).map_err(BacnetError::Codec)?;
if let Some(npdu) = message.npdu {
if let Some(dst_src) = npdu.src {
if let NetworkMessage::Apdu(ApplicationPdu::UnconfirmedRequest(
UnconfirmedRequest::IAm(iam),
)) = npdu.network_message
{
let mut iam = iam.clone();
iam.dst_addr = Some(dst_src);
iams.push(iam);
}
}
}
},
Err(_) => {
break;
}
}
}
Ok(Some(iams))
}
#[maybe_async()]
#[cfg_attr(feature = "alloc", bacnet_macros::remove_lifetimes_from_fn_args)]
pub async fn read_property_multiple<'a>(
&self,
buf: &'a mut [u8],
request: ReadPropertyMultiple<'_>,
addr: Option<NetworkAddress>,
) -> Result<ReadPropertyMultipleAck<'a>, BacnetError<T>> {
let service = ConfirmedRequestService::ReadPropertyMultiple(request);
if let Some(ack) = self.send_and_receive_complex_ack(buf, service, addr).await? {
match ack.service {
ComplexAckService::ReadPropertyMultiple(ack) => Ok(ack),
_ => Err(BacnetError::Codec(Error::ConvertDataLink(
"apdu message is not a ComplexAckService ReadPropertyMultipleAck",
))),
}
} else {
Err(BacnetError::Codec(Error::ConvertDataLink(
"apdu message is not a ComplexAckService ReadPropertyMultipleAck",
)))
}
}
#[maybe_async()]
#[cfg_attr(feature = "alloc", bacnet_macros::remove_lifetimes_from_fn_args)]
pub async fn read_property<'a>(
&self,
buf: &'a mut [u8],
request: ReadProperty,
addr: Option<NetworkAddress>,
) -> Result<ReadPropertyAck<'a>, BacnetError<T>> {
let service = ConfirmedRequestService::ReadProperty(request);
if let Some(ack) = self.send_and_receive_complex_ack(buf, service, addr).await? {
match ack.service {
ComplexAckService::ReadProperty(ack) => Ok(ack),
_ => Err(BacnetError::Codec(Error::ConvertDataLink(
"apdu message is not a ComplexAckService ReadPropertyAck",
))),
}
} else {
Err(BacnetError::Codec(Error::ConvertDataLink(
"apdu message is not a ComplexAckService ReadPropertyAck",
)))
}
}
#[maybe_async()]
pub async fn subscribe_change_of_value(
&self,
buf: &mut [u8],
request: SubscribeCov,
) -> Result<(), BacnetError<T>> {
let service = ConfirmedRequestService::SubscribeCov(request);
let _ack = self.send_and_receive_simple_ack(buf, service, None).await?;
Ok(())
}
#[maybe_async()]
#[cfg_attr(feature = "alloc", bacnet_macros::remove_lifetimes_from_fn_args)]
pub async fn read_change_of_value<'a>(
&self,
buf: &'a mut [u8],
) -> Result<Option<CovNotification<'a>>, BacnetError<T>> {
let n = self.io.read(buf).await.map_err(BacnetError::Io)?;
let mut reader = Reader::default();
let message = DataLink::decode(&mut reader, &buf[..n])?;
if let Some(npdu) = message.npdu {
if let NetworkMessage::Apdu(ApplicationPdu::UnconfirmedRequest(
UnconfirmedRequest::CovNotification(x),
)) = npdu.network_message
{
return Ok(Some(x));
}
};
Ok(None)
}
#[maybe_async()]
#[cfg_attr(feature = "alloc", bacnet_macros::remove_lifetimes_from_fn_args)]
pub async fn read_range<'a>(
&self,
buf: &'a mut [u8],
request: ReadRange,
) -> Result<ReadRangeAck<'a>, BacnetError<T>> {
let service = ConfirmedRequestService::ReadRange(request);
if let Some(ack) = self.send_and_receive_complex_ack(buf, service, None).await? {
match ack.service {
ComplexAckService::ReadRange(ack) => Ok(ack),
_ => Err(BacnetError::Codec(Error::ConvertDataLink(
"apdu message is not a ComplexAckService ReadRangeAck",
))),
}
} else {
Err(BacnetError::Codec(Error::ConvertDataLink(
"apdu message is not a ComplexAckService ReadRangeAck",
)))
}
}
#[maybe_async()]
pub async fn write_property<'a>(
&self,
buf: &mut [u8],
request: WriteProperty<'_>,
addr: Option<NetworkAddress>,
) -> Result<(), BacnetError<T>> {
let service = ConfirmedRequestService::WriteProperty(request);
let _ack = self.send_and_receive_simple_ack(buf, service, addr).await?;
Ok(())
}
#[maybe_async()]
pub async fn write_property_multiple<'a>(
&self,
buf: &mut [u8],
request: WritePropertyMultiple<'_>,
addr: Option<NetworkAddress>,
) -> Result<(), BacnetError<T>> {
let service = ConfirmedRequestService::WritePropertyMultiple(request);
let _ack = self.send_and_receive_simple_ack(buf, service, addr).await?;
Ok(())
}
#[maybe_async()]
pub async fn time_sync(
&self,
buf: &mut [u8],
request: TimeSynchronization,
) -> Result<(), BacnetError<T>> {
let service = UnconfirmedRequest::TimeSynchronization(request);
self.send_unconfirmed(buf, service).await
}
#[maybe_async()]
#[cfg_attr(feature = "alloc", bacnet_macros::remove_lifetimes_from_fn_args)]
async fn send_and_receive_complex_ack<'a>(
&self,
buf: &'a mut [u8],
service: ConfirmedRequestService<'_>,
addr: Option<NetworkAddress>,
) -> Result<Option<ComplexAck<'a>>, BacnetError<T>> {
let invoke_id = self.send_confirmed(buf, service, addr).await?;
for _ in 1..3 {
let n = self.io.read(buf).await.map_err(BacnetError::Io)?;
let buf = &buf[..n];
let mut reader = Reader::default();
let message = DataLink::decode(&mut reader, buf).map_err(BacnetError::Codec)?;
match message.npdu {
Some(x) => match x.network_message {
NetworkMessage::Apdu(ApplicationPdu::ComplexAck(ack)) => {
if ack.invoke_id < invoke_id {
continue;
}
Self::check_invoke_id(invoke_id, ack.invoke_id)?;
return Ok(Some(ack));
}
_ => continue,
},
_ => continue,
}
}
Ok(None)
}
#[maybe_async()]
async fn send_and_receive_simple_ack<'a>(
&self,
buf: &mut [u8],
service: ConfirmedRequestService<'_>,
addr: Option<NetworkAddress>,
) -> Result<SimpleAck, BacnetError<T>> {
let invoke_id = self.send_confirmed(buf, service, addr).await?;
let n = self.io.read(buf).await.map_err(BacnetError::Io)?;
let buf = &buf[..n];
let mut reader = Reader::default();
let message = DataLink::decode(&mut reader, buf).map_err(BacnetError::Codec)?;
let ack: SimpleAck = message.try_into().map_err(BacnetError::Codec)?;
Self::check_invoke_id(invoke_id, ack.invoke_id)?;
Ok(ack)
}
#[maybe_async()]
async fn send_unconfirmed(
&self,
buf: &mut [u8],
service: UnconfirmedRequest<'_>,
) -> Result<(), BacnetError<T>> {
let apdu = ApplicationPdu::UnconfirmedRequest(service);
let message = NetworkMessage::Apdu(apdu);
let npdu = NetworkPdu::new(None, None, true, MessagePriority::Normal, message);
let data_link = DataLink::new(DataLinkFunction::OriginalUnicastNpdu, Some(npdu));
let mut writer = Writer::new(buf);
data_link.encode(&mut writer);
let buffer = writer.to_bytes();
self.io.write(buffer).await.map_err(BacnetError::Io)?;
Ok(())
}
#[maybe_async()]
async fn send_confirmed(
&self,
buf: &mut [u8],
service: ConfirmedRequestService<'_>,
addr: Option<NetworkAddress>,
) -> Result<u8, BacnetError<T>> {
let mut dst: Option<DestinationAddress> = None;
if let Some(addr) = addr {
dst = Some(DestinationAddress::new(addr.net, addr.addr));
}
let invoke_id = self.get_then_inc_invoke_id();
let apdu = ApplicationPdu::ConfirmedRequest(ConfirmedRequest::new(invoke_id, service));
let message = NetworkMessage::Apdu(apdu);
let npdu = NetworkPdu::new(None, dst, true, MessagePriority::Normal, message);
let data_link = DataLink::new(DataLinkFunction::OriginalUnicastNpdu, Some(npdu));
let mut writer = Writer::new(buf);
data_link.encode(&mut writer);
let buffer = writer.to_bytes();
self.io.write(buffer).await.map_err(BacnetError::Io)?;
Ok(invoke_id)
}
fn check_invoke_id(expected: u8, actual: u8) -> Result<(), BacnetError<T>> {
if expected != actual {
Err(BacnetError::InvokeId(InvokeIdError { expected, actual }))
} else {
Ok(())
}
}
fn get_then_inc_invoke_id(&self) -> u8 {
self.invoke_id.fetch_add(1, Ordering::SeqCst)
}
}