use crate::charging_status::*;
use crate::connection_config::*;
use crate::device::*;
use serde_json;
use socket2::{Domain, Socket, Type};
use std::fmt;
use std::net::{Ipv4Addr, UdpSocket};
use std::ops::Drop;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
#[derive(Clone)]
pub struct NetworkAnnouncementMessage {
pub device_name: String,
pub serial_number: String,
pub ip_address: Ipv4Addr,
pub tcp_port: u16,
pub udp_send: u16,
pub udp_receive: u16,
pub rssi: i32,
pub battery: i32,
pub charging_status: ChargingStatus,
pub(crate) time: std::time::Instant,
}
impl From<&NetworkAnnouncementMessage> for Option<TcpConnectionConfig> {
fn from(message: &NetworkAnnouncementMessage) -> Self {
if message.ip_address == Ipv4Addr::UNSPECIFIED || message.tcp_port == 0 {
return None;
}
Some(TcpConnectionConfig {
ip_address: message.ip_address,
port: message.tcp_port,
})
}
}
impl From<&NetworkAnnouncementMessage> for Option<UdpConnectionConfig> {
fn from(message: &NetworkAnnouncementMessage) -> Self {
if message.ip_address == Ipv4Addr::UNSPECIFIED || message.udp_send == 0 || message.udp_receive == 0 {
return None;
}
Some(UdpConnectionConfig {
ip_address: message.ip_address,
send_port: message.udp_receive, receive_port: message.udp_send,
})
}
}
impl From<&NetworkAnnouncementMessage> for Vec<Device> {
fn from(message: &NetworkAnnouncementMessage) -> Self {
let device = |config: Option<ConnectionConfig>| {
Some(Device {
device_name: message.device_name.clone(),
serial_number: message.serial_number.clone(),
connection_config: config?,
})
};
let tcp_device = device(Option::<TcpConnectionConfig>::from(message).map(ConnectionConfig::TcpConnectionConfig));
let udp_device = device(Option::<UdpConnectionConfig>::from(message).map(ConnectionConfig::UdpConnectionConfig));
IntoIterator::into_iter([tcp_device, udp_device]).flatten().collect()
}
}
impl fmt::Display for NetworkAnnouncementMessage {
#[rustfmt::skip]
fn fmt(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
write!(
formatter,
"{}, {}, {}, {}, {}, {}, {}%, {}%, {}",
self.device_name,
self.serial_number,
self.ip_address,
self.tcp_port,
self.udp_send,
self.udp_receive,
self.rssi,
self.battery,
self.charging_status,
)
}
}
impl PartialEq for NetworkAnnouncementMessage {
#[rustfmt::skip]
fn eq(&self, other: &Self) -> bool { self.device_name == other.device_name
&& self.serial_number == other.serial_number
&& self.ip_address == other.ip_address
&& self.tcp_port == other.tcp_port
&& self.udp_send == other.udp_send
&& self.udp_receive == other.udp_receive
}
}
pub struct NetworkAnnouncement {
dropped: Arc<Mutex<bool>>,
closure_counter: AtomicU64,
closures: Arc<Mutex<Vec<(Box<dyn Fn(NetworkAnnouncementMessage) + Send>, u64)>>>,
messages: Arc<Mutex<Vec<NetworkAnnouncementMessage>>>,
}
impl NetworkAnnouncement {
pub fn new() -> std::io::Result<Self> {
let socket: UdpSocket = {
let socket = Socket::new(Domain::IPV4, Type::DGRAM, None)?;
socket.set_reuse_address(true)?;
#[cfg(unix)]
socket.set_reuse_port(true)?;
socket.bind(&"0.0.0.0:10000".parse::<std::net::SocketAddr>().unwrap().into())?;
socket.into()
};
socket.set_nonblocking(true)?;
let network_announcement = Self {
dropped: Arc::new(Mutex::new(false)),
closure_counter: AtomicU64::new(0),
closures: Arc::new(Mutex::new(Vec::new())),
messages: Arc::new(Mutex::new(Vec::new())),
};
let dropped = network_announcement.dropped.clone();
let closures = network_announcement.closures.clone();
let messages = network_announcement.messages.clone();
std::thread::spawn(move || {
loop {
let mut buffer = [0_u8; 1024];
let message = {
match socket.recv_from(&mut buffer) {
Ok((number_of_bytes, _)) => Self::parse(&buffer[..number_of_bytes]),
Err(_) => {
std::thread::sleep(std::time::Duration::from_millis(1));
None
}
}
};
if let Some(message) = &message {
if let Ok(mut messages) = messages.lock() {
if let Some(index) = messages.iter().position(|element| element == message) {
messages[index] = message.clone(); } else {
messages.push(message.clone());
}
}
}
messages.lock().unwrap().retain(|device| device.time.elapsed() < std::time::Duration::from_secs(2));
if let Ok(dropped) = dropped.lock() {
if *dropped {
return;
}
if let Some(message) = message {
closures.lock().unwrap().iter().for_each(|(closure, _)| closure(message.clone()));
}
}
}
});
Ok(network_announcement)
}
fn parse(json: &[u8]) -> Option<NetworkAnnouncementMessage> {
let object: serde_json::Value = serde_json::from_slice(json).ok()?;
Some(NetworkAnnouncementMessage {
device_name: object.get("name").and_then(|value| value.as_str()).unwrap_or("").to_string(),
serial_number: object.get("sn").and_then(|value| value.as_str()).unwrap_or("").to_string(),
ip_address: object.get("ip").and_then(|value| value.as_str()).and_then(|s| s.parse::<Ipv4Addr>().ok()).unwrap_or(Ipv4Addr::UNSPECIFIED),
tcp_port: object.get("port").and_then(|value| value.as_u64()).unwrap_or(0) as u16,
udp_send: object.get("send").and_then(|value| value.as_u64()).unwrap_or(0) as u16,
udp_receive: object.get("receive").and_then(|value| value.as_u64()).unwrap_or(0) as u16,
rssi: object.get("rssi").and_then(|value| value.as_i64()).unwrap_or(-1) as i32,
battery: object.get("battery").and_then(|value| value.as_i64()).unwrap_or(-1) as i32,
charging_status: ChargingStatus::from(object.get("status").and_then(|value| value.as_i64()).unwrap_or(-1) as i32),
time: std::time::Instant::now(),
})
}
pub fn add_closure(&self, closure: Box<dyn Fn(NetworkAnnouncementMessage) + Send>) -> u64 {
let id = self.closure_counter.fetch_add(1, Ordering::SeqCst);
self.closures.lock().unwrap().push((closure, id));
id
}
pub fn remove_closure(&self, id: u64) {
self.closures.lock().unwrap().retain(|(_, closure_id)| closure_id != &id);
}
pub fn get_messages(&self) -> Vec<NetworkAnnouncementMessage> {
(*self.messages.lock().unwrap()).clone()
}
pub fn get_messages_after_short_delay(&self) -> Vec<NetworkAnnouncementMessage> {
std::thread::sleep(std::time::Duration::from_secs(2));
self.get_messages()
}
}
impl Drop for NetworkAnnouncement {
fn drop(&mut self) {
*self.dropped.lock().unwrap() = true;
}
}