use std::{
collections::HashMap,
io::{
Error,
ErrorKind,
},
net::SocketAddr,
ops::{
Deref,
DerefMut,
},
sync::{
atomic::Ordering::Relaxed,
Arc,
Mutex,
mpsc::{
channel,
Receiver,
Sender,
},
},
thread::{
self,
JoinHandle,
},
time::Duration,
};
use atomic::Atomic;
use bytemuck::NoUninit;
use oneshot::Sender as SendOnce;
use crate::{
PresentationType,
primitive,
};
pub use crate::primitive::ConnectionMode;
pub struct Client {
parameter_settings: ParameterSettings,
primitive_client: Arc<primitive::Client>,
selection_state: Atomic<SelectionState>,
selection_mutex: Mutex<()>,
outbox: Mutex<HashMap<u32, (MessageID, SendOnce<Option<Message>>)>>,
system: Mutex<u32>,
}
impl Client {
pub fn new(
parameter_settings: ParameterSettings
) -> Arc<Self> {
Arc::new(Client {
parameter_settings,
primitive_client: primitive::Client::new(),
selection_state: Default::default(),
selection_mutex: Default::default(),
outbox: Default::default(),
system: Default::default(),
})
}
pub fn connect(
self: &Arc<Self>,
entity: &str,
) -> Result<(SocketAddr, Receiver<(MessageID, semi_e5::Message)>), Error> {
let (socket, rx_receiver) = self.primitive_client.connect(entity, self.parameter_settings.connect_mode, self.parameter_settings.t5, self.parameter_settings.t8)?;
let (data_sender, data_receiver) = channel::<(MessageID, semi_e5::Message)>();
let clone: Arc<Client> = self.clone();
thread::spawn(move || {clone.receive(rx_receiver, data_sender)});
Ok((socket, data_receiver))
}
pub fn disconnect(
self: &Arc<Self>,
) -> Result<(), Error> {
let result: Result<(), Error> = self.primitive_client.disconnect();
let _guard = self.selection_mutex.lock().unwrap();
if let SelectionState::Selected = self.selection_state.load(Relaxed) {
self.selection_state.store(SelectionState::NotSelected, Relaxed);
}
result
}
}
impl Client {
fn receive(
self: &Arc<Self>,
rx_receiver: Receiver<primitive::Message>,
rx_sender: Sender<(MessageID, semi_e5::Message)>,
) {
for primitive_message in rx_receiver {
let primitive_header = primitive_message.header;
match Message::try_from(primitive_message) {
Ok(rx_message) => match rx_message.contents {
MessageContents::DataMessage(data) => {
match self.selection_state.load(Relaxed) {
SelectionState::Selected => {
if data.function % 2 == 1 {
if rx_sender.send((rx_message.id, data)).is_err() {break}
}
else {
let mut outbox = self.outbox.lock().unwrap();
let mut optional_transaction: Option<u32> = None;
for (outbox_id, (message_id, _)) in outbox.deref() {
if *message_id == rx_message.id {
optional_transaction = Some(*outbox_id);
break;
}
}
if let Some(transaction) = optional_transaction {
let (_, sender) = outbox.deref_mut().remove(&transaction).unwrap();
sender.send(Some(Message{
id: rx_message.id,
contents: MessageContents::DataMessage(data),
})).unwrap();
}
else {
if self.primitive_client.transmit(Message {
id: rx_message.id,
contents: MessageContents::RejectRequest(0, RejectReason::TransactionNotOpen as u8)
}.into()).is_err() {break}
}
}
},
_ => {
if self.primitive_client.transmit(Message {
id: rx_message.id,
contents: MessageContents::RejectRequest(0, RejectReason::EntityNotSelected as u8)
}.into()).is_err() {break}
},
}
},
MessageContents::SelectRequest => {
match self.selection_mutex.try_lock() {
Ok(_guard) => {
match self.selection_state.load(Relaxed) {
SelectionState::NotSelected => {
if self.primitive_client.transmit(Message {
id: rx_message.id,
contents: MessageContents::SelectResponse(SelectStatus::Success as u8),
}.into()).is_err() {break};
self.selection_state.store(SelectionState::Selected, Relaxed);
},
SelectionState::Selected => {
if self.primitive_client.transmit(Message {
id: rx_message.id,
contents: MessageContents::SelectResponse(SelectStatus::AlreadyActive as u8),
}.into()).is_err() {break};
},
}
},
Err(_) => {
},
}
},
MessageContents::SelectResponse(select_status) => {
let mut outbox = self.outbox.lock().unwrap();
let mut optional_transaction: Option<u32> = None;
for (outbox_id, (message_id, _)) in outbox.deref() {
if *message_id == rx_message.id {
optional_transaction = Some(*outbox_id);
break;
}
}
if let Some(transaction) = optional_transaction {
let (_, sender) = outbox.deref_mut().remove(&transaction).unwrap();
sender.send(Some(Message{
id: rx_message.id,
contents: MessageContents::SelectResponse(select_status),
})).unwrap();
}
else {
if self.primitive_client.transmit(Message {
id: rx_message.id,
contents: MessageContents::RejectRequest(0, RejectReason::TransactionNotOpen as u8)
}.into()).is_err() {break}
}
},
MessageContents::DeselectRequest => {
todo!()
},
MessageContents::DeselectResponse(_deselect_status) => {
todo!()
},
MessageContents::LinktestRequest => {
if self.primitive_client.transmit(Message{
id: rx_message.id,
contents: MessageContents::LinktestResponse,
}.into()).is_err() {break};
},
MessageContents::LinktestResponse => {
let mut outbox = self.outbox.lock().unwrap();
let mut optional_transaction: Option<u32> = None;
for (outbox_id, (message_id, _)) in outbox.deref() {
if *message_id == rx_message.id {
optional_transaction = Some(*outbox_id);
break;
}
}
if let Some(transaction) = optional_transaction {
let (_, sender) = outbox.deref_mut().remove(&transaction).unwrap();
sender.send(Some(rx_message)).unwrap();
}
else {
if self.primitive_client.transmit(Message {
id: rx_message.id,
contents: MessageContents::RejectRequest(SessionType::LinktestRequest as u8, RejectReason::TransactionNotOpen as u8),
}.into()).is_err() {break}
}
},
MessageContents::RejectRequest(_message_type, _reason_code) => {
let mut outbox = self.outbox.lock().unwrap();
let mut optional_transaction: Option<u32> = None;
for (outbox_id, (message_id, _)) in outbox.deref() {
if *message_id == rx_message.id {
optional_transaction = Some(*outbox_id);
break;
}
}
if let Some(transaction) = optional_transaction {
let (_, sender) = outbox.deref_mut().remove(&transaction).unwrap();
sender.send(None).unwrap();
}
},
MessageContents::SeparateRequest => {
let _guard: std::sync::MutexGuard<'_, ()> = self.selection_mutex.lock().unwrap();
if let SelectionState::Selected = self.selection_state.load(Relaxed) {
self.selection_state.store(SelectionState::NotSelected, Relaxed);
}
},
},
Err(reject_reason) => {
if self.primitive_client.transmit(Message {
id: MessageID {
session: primitive_header.session_id,
system: primitive_header.system,
},
contents: MessageContents::RejectRequest(match reject_reason {
RejectReason::UnsupportedPresentationType => primitive_header.presentation_type,
_ => primitive_header.session_type,
}, reject_reason as u8),
}.into()).is_err() {break}
},
}
}
for (_, (_, sender)) in self.outbox.lock().unwrap().deref_mut().drain() {
let _ = sender.send(None);
}
}
fn transmit(
self: &Arc<Self>,
message: Message,
reply_expected: bool,
delay: Duration,
) -> Result<Option<Message>, Error> {
let (receiver, system) = {
let outbox_lock = if reply_expected {Some(self.deref().outbox.lock().unwrap())} else {None};
let message_id = message.id;
match self.primitive_client.transmit(message.into()) {
Ok(()) => {
match outbox_lock {
None => return Ok(None),
Some(mut outbox) => {
let (sender, receiver) = oneshot::channel::<Option<Message>>();
let system = {
let mut system_guard = self.deref().system.lock().unwrap();
let system_counter = system_guard.deref_mut();
let system = *system_counter;
*system_counter += 1;
system
};
outbox.deref_mut().insert(system, (message_id, sender));
(receiver, system)
}
}
},
Err(error) => {
let _ = self.disconnect();
return Err(error)
},
}
};
let rx_result = receiver.recv_timeout(delay);
let mut outbox = self.outbox.lock().unwrap();
outbox.deref_mut().remove(&system);
match rx_result {
Ok(rx_message) => return Ok(rx_message),
Err(_e) => return Ok(None),
}
}
pub fn data(
self: &Arc<Self>,
id: MessageID,
message: semi_e5::Message,
) -> JoinHandle<Result<Option<semi_e5::Message>, Error>> {
let clone: Arc<Client> = self.clone();
let reply_expected: bool = message.function % 2 == 1 && message.w;
thread::spawn(move || {
match clone.selection_state.load(Relaxed) {
SelectionState::NotSelected => return Err(Error::from(ErrorKind::AlreadyExists)),
SelectionState::Selected => {
match clone.transmit(
Message {
id,
contents: MessageContents::DataMessage(message),
},
reply_expected,
clone.parameter_settings.t3,
)?{
Some(rx_message) => {
match rx_message.contents {
MessageContents::DataMessage(data_message) => return Ok(Some(data_message)),
MessageContents::RejectRequest(_type, _reason) => return Err(Error::from(ErrorKind::PermissionDenied)),
_ => return Err(Error::from(ErrorKind::InvalidData)),
}
},
None => {
if reply_expected {
clone.disconnect()?;
Err(Error::from(ErrorKind::ConnectionAborted))
}
else {
return Ok(None);
}
},
}
},
}
})
}
pub fn select(
self: &Arc<Self>,
id: MessageID,
) -> JoinHandle<Result<(), Error>> {
let clone: Arc<Client> = self.clone();
thread::spawn(move || {
'disconnect: {
let _guard = clone.selection_mutex.lock();
match clone.selection_state.load(Relaxed) {
SelectionState::NotSelected => {
match clone.transmit(
Message {
id,
contents: MessageContents::SelectRequest,
},
true,
clone.parameter_settings.t6,
)?{
Some(rx_message) => {
match rx_message.contents {
MessageContents::SelectResponse(select_status) => {
if select_status == SelectStatus::Success as u8 {
clone.selection_state.store(SelectionState::Selected, Relaxed);
return Ok(())
}
else {
return Err(Error::from(ErrorKind::PermissionDenied))
}
},
MessageContents::RejectRequest(_type, _reason) => return Err(Error::from(ErrorKind::PermissionDenied)),
_ => return Err(Error::from(ErrorKind::InvalidData)),
}
},
None => {
break 'disconnect;
},
}
},
SelectionState::Selected => {
return Err(Error::from(ErrorKind::AlreadyExists))
},
}
}
clone.disconnect()?;
Err(Error::from(ErrorKind::ConnectionAborted))
})
}
pub fn deselect(
self: &Arc<Self>,
) -> Result<(), Error> {
todo!()
}
pub fn linktest(
self: &Arc<Self>,
system: u32,
) -> JoinHandle<Result<(), Error>> {
let clone: Arc<Client> = self.clone();
thread::spawn(move || {
match clone.transmit(
Message {
id: MessageID {
session: 0xFFFF,
system,
},
contents: MessageContents::LinktestRequest,
},
true,
clone.parameter_settings.t6,
)?{
Some(rx_message) => {
match rx_message.contents {
MessageContents::LinktestResponse => Ok(()),
MessageContents::RejectRequest(_type, _reason) => Err(Error::from(ErrorKind::PermissionDenied)),
_ => Err(Error::from(ErrorKind::InvalidData)),
}
},
None => {
clone.disconnect()?;
Err(Error::from(ErrorKind::ConnectionAborted))
},
}
})
}
pub fn separate(
self: &Arc<Self>,
id: MessageID,
) -> JoinHandle<Result<(), Error>> {
let clone: Arc<Client> = self.clone();
thread::spawn(move || {
let _guard = clone.selection_mutex.lock().unwrap();
match clone.selection_state.load(Relaxed) {
SelectionState::NotSelected => {
Err(Error::from(ErrorKind::PermissionDenied))
},
SelectionState::Selected => {
clone.transmit(
Message {
id,
contents: MessageContents::SeparateRequest,
},
false,
clone.parameter_settings.t6,
)?;
clone.selection_state.store(SelectionState::NotSelected, Relaxed);
Ok(())
},
}
})
}
pub fn reject(
self: &Arc<Self>,
_reason: RejectReason,
) -> Result<(), Error> {
todo!()
}
}
#[derive(Clone, Copy, Debug, PartialEq, NoUninit)]
#[repr(u8)]
pub enum SelectionState {
NotSelected,
Selected,
}
impl Default for SelectionState {
fn default() -> Self {
SelectionState::NotSelected
}
}
#[derive(Clone, Copy, Debug, PartialEq)]
pub struct ParameterSettings {
pub connect_mode: ConnectionMode,
pub t3: Duration,
pub t5: Duration,
pub t6: Duration,
pub t7: Duration,
pub t8: Duration,
}
impl Default for ParameterSettings {
fn default() -> Self {
Self {
connect_mode: ConnectionMode::default(),
t3: Duration::from_secs(45),
t5: Duration::from_secs(10),
t6: Duration::from_secs(5),
t7: Duration::from_secs(10),
t8: Duration::from_secs(5),
}
}
}
#[derive(Clone, Debug)]
pub struct Message {
pub id: MessageID,
pub contents: MessageContents,
}
impl From<Message> for primitive::Message {
fn from(message: Message) -> Self {
match message.contents {
MessageContents::DataMessage(e5_message) => {
primitive::Message {
header: primitive::MessageHeader {
session_id : message.id.session,
byte_2 : ((e5_message.w as u8) << 7) | e5_message.stream,
byte_3 : e5_message.function,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::DataMessage as u8,
system : message.id.system,
},
text: match e5_message.text {
Some(item) => Vec::<u8>::from(item),
None => vec![],
},
}
},
MessageContents::SelectRequest => {
primitive::Message {
header: primitive::MessageHeader {
session_id : message.id.session,
byte_2 : 0,
byte_3 : 0,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::SelectRequest as u8,
system : message.id.system,
},
text: vec![],
}
},
MessageContents::SelectResponse(select_status) => {
primitive::Message {
header: primitive::MessageHeader {
session_id : message.id.session,
byte_2 : 0,
byte_3 : select_status,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::SelectResponse as u8,
system : message.id.system,
},
text: vec![],
}
},
MessageContents::DeselectRequest => {
primitive::Message {
header: primitive::MessageHeader {
session_id : message.id.session,
byte_2 : 0,
byte_3 : 0,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::DeselectRequest as u8,
system : message.id.system,
},
text: vec![],
}
},
MessageContents::DeselectResponse(deselect_status) => {
primitive::Message {
header: primitive::MessageHeader {
session_id : message.id.session,
byte_2 : 0,
byte_3 : deselect_status,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::DeselectResponse as u8,
system : message.id.system,
},
text: vec![],
}
},
MessageContents::LinktestRequest => {
primitive::Message {
header: primitive::MessageHeader {
session_id : 0xFFFF,
byte_2 : 0,
byte_3 : 0,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::LinktestRequest as u8,
system : message.id.system,
},
text: vec![],
}
},
MessageContents::LinktestResponse => {
primitive::Message {
header: primitive::MessageHeader {
session_id : 0xFFFF,
byte_2 : 0,
byte_3 : 0,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::LinktestResponse as u8,
system : message.id.system,
},
text: vec![],
}
},
MessageContents::RejectRequest(message_type, reason_code) => {
primitive::Message {
header: primitive::MessageHeader {
session_id : message.id.session,
byte_2 : message_type,
byte_3 : reason_code,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::RejectRequest as u8,
system : message.id.system,
},
text: vec![],
}
},
MessageContents::SeparateRequest => {
primitive::Message {
header: primitive::MessageHeader {
session_id : message.id.session,
byte_2 : 0,
byte_3 : 0,
presentation_type : PresentationType::SecsII as u8,
session_type : SessionType::SeparateRequest as u8,
system : message.id.system,
},
text: vec![],
}
},
}
}
}
impl TryFrom<primitive::Message> for Message {
type Error = RejectReason;
fn try_from(message: primitive::Message) -> Result<Self, Self::Error> {
if message.header.presentation_type != 0 {return Err(RejectReason::UnsupportedPresentationType)}
Ok(Message {
id: MessageID {
session: message.header.session_id,
system: message.header.system,
},
contents: match message.header.session_type {
0 => {
MessageContents::DataMessage(semi_e5::Message{
stream : message.header.byte_2 & 0b0111_1111,
function : message.header.byte_3,
w : message.header.byte_2 & 0b1000_0000 > 0,
text : match semi_e5::Item::try_from(message.text) {
Ok(text) => Some(text),
Err(error) => {
match error {
semi_e5::Error::EmptyText => {None},
_ => {return Err(RejectReason::MalformedData)}
}
},
},
})
},
1 => {
if message.header.byte_2 != 0 {return Err(RejectReason::MalformedData)}
if message.header.byte_3 != 0 {return Err(RejectReason::MalformedData)}
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::SelectRequest
},
2 => {
if message.header.byte_2 != 0 {return Err(RejectReason::MalformedData)}
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::SelectResponse(message.header.byte_3)
},
3 => {
if message.header.byte_2 != 0 {return Err(RejectReason::MalformedData)}
if message.header.byte_3 != 0 {return Err(RejectReason::MalformedData)}
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::DeselectRequest
},
4 => {
if message.header.byte_2 != 0 {return Err(RejectReason::MalformedData)}
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::DeselectResponse(message.header.byte_3)
},
5 => {
if message.header.session_id != 0xFFFF {return Err(RejectReason::MalformedData)}
if message.header.byte_2 != 0 {return Err(RejectReason::MalformedData)}
if message.header.byte_3 != 0 {return Err(RejectReason::MalformedData)}
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::LinktestRequest
},
6 => {
if message.header.session_id != 0xFFFF {return Err(RejectReason::MalformedData)}
if message.header.byte_2 != 0 {return Err(RejectReason::MalformedData)}
if message.header.byte_3 != 0 {return Err(RejectReason::MalformedData)}
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::LinktestResponse
},
7 => {
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::RejectRequest(message.header.byte_2, message.header.byte_3)
},
9 => {
if message.header.byte_2 != 0 {return Err(RejectReason::MalformedData)}
if message.header.byte_3 != 0 {return Err(RejectReason::MalformedData)}
if !message.text.is_empty() {return Err(RejectReason::MalformedData)}
MessageContents::SeparateRequest
},
_ => {return Err(RejectReason::UnsupportedSessionType)}
},
})
}
}
#[derive(Clone, Copy, Debug, PartialEq)]
pub struct MessageID {
pub session: u16,
pub system: u32,
}
#[repr(u8)]
#[derive(Clone, Debug)]
pub enum MessageContents {
DataMessage(semi_e5::Message) = SessionType::DataMessage as u8,
SelectRequest = SessionType::SelectRequest as u8,
SelectResponse(u8) = SessionType::SelectResponse as u8,
DeselectRequest = SessionType::DeselectRequest as u8,
DeselectResponse(u8) = SessionType::DeselectResponse as u8,
LinktestRequest = SessionType::LinktestRequest as u8,
LinktestResponse = SessionType::LinktestResponse as u8,
RejectRequest(u8, u8) = SessionType::RejectRequest as u8,
SeparateRequest = SessionType::SeparateRequest as u8,
}
#[repr(u8)]
#[derive(Clone, Copy, Debug, PartialEq)]
pub enum SessionType {
DataMessage = 0,
SelectRequest = 1,
SelectResponse = 2,
DeselectRequest = 3,
DeselectResponse = 4,
LinktestRequest = 5,
LinktestResponse = 6,
RejectRequest = 7,
SeparateRequest = 9,
}
#[repr(u8)]
#[derive(Clone, Copy, Debug, PartialEq)]
pub enum SelectStatus {
Success = 0,
AlreadyActive = 1,
NotReady = 2,
Exhausted = 3,
}
#[repr(u8)]
#[derive(Clone, Copy, Debug, PartialEq)]
pub enum DeselectStatus {
Success = 0,
NotEstablished = 1,
Busy = 2,
}
#[repr(u8)]
#[derive(Clone, Copy, Debug, PartialEq)]
pub enum RejectReason {
MalformedData = 0,
UnsupportedSessionType = 1,
UnsupportedPresentationType = 2,
TransactionNotOpen = 3,
EntityNotSelected = 4,
}