From 96f7888b08fa9eb9cb140b330a33b5c151d1e4bb Mon Sep 17 00:00:00 2001 From: "t.galimkhanov" Date: Thu, 7 May 2026 05:38:08 +0500 Subject: [PATCH] improve: async - homework 3 --- .gitignore | 3 +- examples/simulated_home.rs | 182 ++++++++++++++++++ examples/thermo_kitchen.txt | 3 + examples/thermo_living.txt | 3 + src/bin/socket_sim.rs | 187 ++++++++++++++++++ src/bin/thermo_sim.rs | 108 +++++++++++ src/devices/socket.rs | 372 ++++++++++++++++++++++++++++-------- src/devices/thermometer.rs | 247 ++++++++++++++++++++---- src/lib.rs | 2 + src/wire.rs | 151 +++++++++++++++ 10 files changed, 1149 insertions(+), 109 deletions(-) create mode 100644 examples/simulated_home.rs create mode 100644 examples/thermo_kitchen.txt create mode 100644 examples/thermo_living.txt create mode 100644 src/bin/socket_sim.rs create mode 100644 src/bin/thermo_sim.rs create mode 100644 src/wire.rs diff --git a/.gitignore b/.gitignore index 68b6927..eb68687 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,2 @@ /target -/tmp -/examples \ No newline at end of file +/tmp \ No newline at end of file diff --git a/examples/simulated_home.rs b/examples/simulated_home.rs new file mode 100644 index 0000000..3129146 --- /dev/null +++ b/examples/simulated_home.rs @@ -0,0 +1,182 @@ +//! Example smart home wired to `socket_sim` / `thermo_sim` processes. +//! +//! Typical local demo (four terminals): +//! +//! ```text +//! cargo run --bin socket_sim -- 127.0.0.1:17001 --watts 200 +//! cargo run --bin socket_sim -- 127.0.0.1:17002 --watts 40 +//! cargo run --bin thermo_sim -- examples/thermo_kitchen.txt +//! cargo run --bin thermo_sim -- examples/thermo_living.txt +//! cargo run --example simulated_home +//! ``` +//! +//! Environment overrides (optional): +//! - `SIM_SOCKET_FRIDGE`, `SIM_SOCKET_DESK` — TCP addresses for outlets +//! - `SIM_THERMO_KITCHEN_BIND`, `SIM_THERMO_LIVING_BIND` — UDP bind addresses for thermometers + +use std::collections::HashMap; +use std::thread; +use std::time::Duration; + +use smart_home::{Device, Power, SmartHome, Socket, Temperature, Thermometer, print_report_value}; + +fn env_addr(name: &str, default: &str) -> String { + std::env::var(name).unwrap_or_else(|_| default.to_string()) +} + +fn main() { + let fridge_addr = env_addr("SIM_SOCKET_FRIDGE", "127.0.0.1:17001"); + let desk_addr = env_addr("SIM_SOCKET_DESK", "127.0.0.1:17002"); + let kitchen_bind = env_addr("SIM_THERMO_KITCHEN_BIND", "127.0.0.1:19101"); + let living_bind = env_addr("SIM_THERMO_LIVING_BIND", "127.0.0.1:19102"); + + let mut errors: Vec = Vec::new(); + + let fridge_socket = match fridge_addr.parse() { + Ok(addr) => match Socket::connect_tcp( + "Fridge outlet".to_string(), + addr, + Power::new(200.0).unwrap(), + ) { + Ok(s) => Some(s), + Err(e) => { + errors.push(format!("Fridge TCP outlet ({fridge_addr}): {e}")); + None + } + }, + Err(e) => { + errors.push(format!("Fridge address ({fridge_addr}): {e}")); + None + } + }; + + let desk_socket = match desk_addr.parse() { + Ok(addr) => match Socket::connect_tcp( + "Desk lamp outlet".to_string(), + addr, + Power::new(40.0).unwrap(), + ) { + Ok(s) => Some(s), + Err(e) => { + errors.push(format!("Desk TCP outlet ({desk_addr}): {e}")); + None + } + }, + Err(e) => { + errors.push(format!("Desk address ({desk_addr}): {e}")); + None + } + }; + + let kitchen_thermo = match kitchen_bind.parse() { + Ok(addr) => match Thermometer::bind_udp( + "Kitchen sensor".to_string(), + addr, + Temperature::celsius(20.0), + ) { + Ok(t) => Some(t), + Err(e) => { + errors.push(format!( + "Kitchen UDP thermometer bind ({kitchen_bind}): {e}" + )); + None + } + }, + Err(e) => { + errors.push(format!("Kitchen bind address ({kitchen_bind}): {e}")); + None + } + }; + + let living_thermo = match living_bind.parse() { + Ok(addr) => match Thermometer::bind_udp( + "Living room sensor".to_string(), + addr, + Temperature::celsius(20.0), + ) { + Ok(t) => Some(t), + Err(e) => { + errors.push(format!( + "Living room UDP thermometer bind ({living_bind}): {e}" + )); + None + } + }, + Err(e) => { + errors.push(format!("Living room bind address ({living_bind}): {e}")); + None + } + }; + + // Give UDP senders time to deliver at least one datagram. + thread::sleep(Duration::from_millis(600)); + + let kitchen_thermo = kitchen_thermo.and_then(|t| { + if t.is_udp() && !t.has_udp_reading() { + errors.push(format!( + "Kitchen thermometer ({kitchen_bind}): no UDP temperature received yet" + )); + None + } else { + Some(t) + } + }); + let living_thermo = living_thermo.and_then(|t| { + if t.is_udp() && !t.has_udp_reading() { + errors.push(format!( + "Living room thermometer ({living_bind}): no UDP temperature received yet" + )); + None + } else { + Some(t) + } + }); + + let mut kitchen_devices: HashMap = HashMap::new(); + if let Some(t) = kitchen_thermo { + kitchen_devices.insert("kitchen_thermometer".to_string(), t.into()); + } + if let Some(s) = fridge_socket { + kitchen_devices.insert("fridge".to_string(), s.into()); + } + + let mut office_devices: HashMap = HashMap::new(); + if let Some(t) = living_thermo { + office_devices.insert("living_thermometer".to_string(), t.into()); + } + if let Some(s) = desk_socket { + office_devices.insert("desk".to_string(), s.into()); + } + + let mut rooms = HashMap::new(); + if !kitchen_devices.is_empty() { + rooms.insert( + "Kitchen".to_string(), + smart_home::Room::new("Kitchen".to_string(), kitchen_devices), + ); + } + if !office_devices.is_empty() { + rooms.insert( + "Office".to_string(), + smart_home::Room::new("Office".to_string(), office_devices), + ); + } + + let home = SmartHome::new("Simulated smart home".to_string(), rooms); + + println!("=== Simulated smart home report ===\n"); + if home.room_count() == 0 { + println!("No devices were added (all connections/bindings failed)."); + } else { + print_report_value(&home); + } + + if errors.is_empty() { + println!("\nNo transport-level errors while constructing devices."); + } else { + println!("\nTransport / data errors:"); + for e in &errors { + println!("- {e}"); + } + } +} diff --git a/examples/thermo_kitchen.txt b/examples/thermo_kitchen.txt new file mode 100644 index 0000000..ae40a60 --- /dev/null +++ b/examples/thermo_kitchen.txt @@ -0,0 +1,3 @@ +127.0.0.1:19101 +300 +21.0 diff --git a/examples/thermo_living.txt b/examples/thermo_living.txt new file mode 100644 index 0000000..3a65559 --- /dev/null +++ b/examples/thermo_living.txt @@ -0,0 +1,3 @@ +127.0.0.1:19102 +300 +23.5 diff --git a/src/bin/socket_sim.rs b/src/bin/socket_sim.rs new file mode 100644 index 0000000..1b8c078 --- /dev/null +++ b/src/bin/socket_sim.rs @@ -0,0 +1,187 @@ +//! TCP smart-outlet simulator: non-blocking accept/read, shared state, many clients. +//! +//! Usage: `socket_sim [--watts ] [--on|--off]` + +use std::io::{self, Read, Write}; +use std::net::{SocketAddr, TcpListener, TcpStream}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use smart_home::wire; + +#[derive(Debug, Clone)] +struct OutletState { + is_on: bool, + nominal_watts: f32, +} + +impl OutletState { + fn watts_now(&self) -> f32 { + if self.is_on { self.nominal_watts } else { 0.0 } + } +} + +struct Client { + stream: TcpStream, + out: Vec, + buf: Vec, +} + +impl Client { + fn new(stream: TcpStream) -> io::Result { + stream.set_nonblocking(true)?; + Ok(Self { + stream, + out: Vec::new(), + buf: Vec::new(), + }) + } + + fn push_response(&mut self, line: &str) { + self.out.extend_from_slice(line.as_bytes()); + if !line.ends_with('\n') { + self.out.push(b'\n'); + } + } + + fn flush_writes(&mut self) -> io::Result<()> { + while !self.out.is_empty() { + match self.stream.write(&self.out) { + Ok(0) => { + return Err(io::Error::new(io::ErrorKind::WriteZero, "short write")); + } + Ok(n) => { + self.out.drain(..n); + } + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => break, + Err(e) => return Err(e), + } + } + Ok(()) + } + + fn read_available(&mut self) -> io::Result<()> { + let mut tmp = [0_u8; 512]; + loop { + match self.stream.read(&mut tmp) { + Ok(0) => { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "client closed", + )); + } + Ok(n) => self.buf.extend_from_slice(&tmp[..n]), + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => break, + Err(e) => return Err(e), + } + } + Ok(()) + } + + fn drain_lines(&mut self, state: &Arc>) -> io::Result<()> { + while let Some(pos) = self.buf.iter().position(|&b| b == b'\n') { + let raw: Vec = self.buf.drain(..=pos).collect(); + let line = String::from_utf8_lossy(&raw[..raw.len().saturating_sub(1)]); + let line = line.trim(); + match line { + "GET_STATUS" => { + let g = state + .lock() + .map_err(|_| io::Error::other("state mutex poisoned"))?; + let msg = wire::format_status_line(g.is_on, g.watts_now()); + self.push_response(msg.trim_end()); + } + "SET_ON" => { + if let Ok(mut g) = state.lock() { + g.is_on = true; + } + self.push_response("OK"); + } + "SET_OFF" => { + if let Ok(mut g) = state.lock() { + g.is_on = false; + } + self.push_response("OK"); + } + "" => {} + other => { + self.push_response(&format!("ERR unknown command: {other}")); + } + } + } + Ok(()) + } +} + +fn parse_args() -> Result<(SocketAddr, OutletState), String> { + let mut args = std::env::args().skip(1); + let bind: SocketAddr = args + .next() + .ok_or_else(|| "missing BIND_ADDR".to_string())? + .parse() + .map_err(|e| format!("invalid bind address: {e}"))?; + + let mut watts = 1500.0_f32; + let mut is_on = false; + while let Some(a) = args.next() { + match a.as_str() { + "--watts" => { + let v = args + .next() + .ok_or_else(|| "--watts needs a value".to_string())?; + watts = v.parse().map_err(|e| format!("invalid watts: {e}"))?; + } + "--on" => is_on = true, + "--off" => is_on = false, + other => return Err(format!("unknown arg: {other}")), + } + } + + Ok(( + bind, + OutletState { + is_on, + nominal_watts: watts, + }, + )) +} + +fn main() -> Result<(), Box> { + let (bind, initial) = + parse_args().map_err(|e| io::Error::new(io::ErrorKind::InvalidInput, e))?; + let state = Arc::new(Mutex::new(initial)); + let listener = TcpListener::bind(bind)?; + listener.set_nonblocking(true)?; + eprintln!("socket_sim listening on {bind} (non-blocking)"); + + let mut clients: Vec = Vec::new(); + + loop { + match listener.accept() { + Ok((stream, _)) => match Client::new(stream) { + Ok(c) => clients.push(c), + Err(e) => eprintln!("client init: {e}"), + }, + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {} + Err(e) => return Err(e.into()), + } + + let mut i = 0; + while i < clients.len() { + let remove = { + let c = &mut clients[i]; + match c.read_available() { + Err(_) => true, + Ok(()) => c.drain_lines(&state).is_err() || c.flush_writes().is_err(), + } + }; + if remove { + clients.swap_remove(i); + } else { + i += 1; + } + } + + std::thread::sleep(Duration::from_millis(2)); + } +} diff --git a/src/bin/thermo_sim.rs b/src/bin/thermo_sim.rs new file mode 100644 index 0000000..9c460b3 --- /dev/null +++ b/src/bin/thermo_sim.rs @@ -0,0 +1,108 @@ +//! UDP temperature simulator: non-blocking send loop, config file driven. +//! +//! Config file format (UTF-8 text): +//! - line 1: destination `HOST:PORT` for UDP datagrams +//! - line 2: period in milliseconds between sends +//! - line 3 (optional): fixed temperature in °C; if omitted, each send uses a pseudo-random value +//! +//! Usage: `thermo_sim ` + +use std::fs; +use std::io; +use std::net::{SocketAddr, UdpSocket}; +use std::time::{Duration, Instant}; + +use smart_home::wire; + +struct Config { + dest: SocketAddr, + period: Duration, + fixed: Option, +} + +fn parse_config(text: &str) -> Result { + let mut lines = text.lines().map(str::trim).filter(|l| !l.is_empty()); + let dest_line = lines + .next() + .ok_or_else(|| "missing destination line".to_string())?; + let dest: SocketAddr = dest_line + .parse() + .map_err(|e| format!("invalid destination: {e}"))?; + let period_line = lines + .next() + .ok_or_else(|| "missing period line".to_string())?; + let ms: u64 = period_line + .parse() + .map_err(|e| format!("invalid period ms: {e}"))?; + let fixed = if let Some(t) = lines.next() { + Some( + t.parse::() + .map_err(|e| format!("invalid temperature: {e}"))?, + ) + } else { + None + }; + Ok(Config { + dest, + period: Duration::from_millis(ms), + fixed, + }) +} + +fn next_temp(fixed: Option, tick: u64) -> f32 { + if let Some(t) = fixed { + return t; + } + // deterministic "random" in a comfortable band for demos + let x = ((tick.wrapping_mul(6364136223846793005) ^ 0x9E37_79B9) % 10_000) as f32 / 10_000.0; + 18.0 + x * 8.0 +} + +fn main() { + let path = std::env::args().nth(1).unwrap_or_else(|| usage_and_exit()); + let text = fs::read_to_string(&path).unwrap_or_else(|e| { + eprintln!("read config: {e}"); + std::process::exit(1); + }); + let cfg = parse_config(&text).unwrap_or_else(|e| { + eprintln!("config: {e}"); + std::process::exit(1); + }); + + let sock = UdpSocket::bind("0.0.0.0:0").unwrap_or_else(|e| { + eprintln!("bind udp: {e}"); + std::process::exit(1); + }); + sock.set_nonblocking(true).unwrap_or_else(|e| { + eprintln!("set_nonblocking: {e}"); + std::process::exit(1); + }); + + eprintln!( + "thermo_sim -> {} every {:?} (non-blocking)", + cfg.dest, cfg.period + ); + + let mut next = Instant::now(); + let mut tick: u64 = 0; + loop { + let now = Instant::now(); + if now >= next { + let t = next_temp(cfg.fixed, tick); + tick = tick.wrapping_add(1); + let payload = wire::encode_temperature_celsius(t); + match sock.send_to(&payload, cfg.dest) { + Ok(_) => {} + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {} + Err(e) => eprintln!("send: {e}"), + } + next = now.checked_add(cfg.period).unwrap_or(now + cfg.period); + } + std::thread::sleep(Duration::from_millis(1)); + } +} + +fn usage_and_exit() -> ! { + eprintln!("usage: thermo_sim "); + std::process::exit(2); +} diff --git a/src/devices/socket.rs b/src/devices/socket.rs index dbd47fa..6a16fbd 100644 --- a/src/devices/socket.rs +++ b/src/devices/socket.rs @@ -1,113 +1,335 @@ +use std::fmt; +use std::net::SocketAddr; +use std::sync::Mutex; + use super::device::DeviceInfo; use crate::types::Power; +use crate::wire; -/// Smart socket for controlling electrical appliances -#[derive(Debug, Clone, PartialEq)] +/// Smart socket: local state or synchronous TCP to a remote outlet. +#[derive(Debug)] pub struct Socket { name: String, + inner: SocketInner, +} + +#[derive(Debug)] +enum SocketInner { + Local { + is_on: bool, + power: Power, + }, + Tcp { + addr: SocketAddr, + /// Rated power (watts) when the outlet is on; used if the remote does not echo watts. + rated: Power, + state: Mutex, + }, +} + +#[derive(Debug, Clone)] +struct TcpState { is_on: bool, - power: Power, + watts: f32, + last_error: Option, } impl Socket { - /// Creates a new socket - /// - /// # Arguments - /// - /// * `name` - Socket name - /// * `is_on` - Initial state (true - on, false - off) - /// * `power` - Power of the connected device - /// - /// # Examples - /// - /// ``` - /// use smart_home::{Socket, Power, DeviceInfo}; - /// - /// let socket = Socket::new("Kettle".to_string(), true, Power::new(1500.0).unwrap()); - /// assert_eq!(socket.name(), "Kettle"); - /// assert!(socket.is_on()); - /// ``` + /// Creates a new local socket (no network I/O). pub fn new(name: String, is_on: bool, power: Power) -> Self { - Self { name, is_on, power } + Self { + name, + inner: SocketInner::Local { is_on, power }, + } } - /// Turn on the socket - /// - /// # Examples + /// Connects to a remote smart outlet over TCP and reads its current state. /// - /// ``` - /// use smart_home::{Socket, Power}; - /// - /// let mut socket = Socket::new("Lamp".to_string(), false, Power::new(60.0).unwrap()); - /// socket.turn_on(); - /// assert!(socket.is_on()); - /// ``` + /// This performs a synchronous `GET_STATUS` round-trip during construction. + pub fn connect_tcp(name: String, addr: SocketAddr, rated: Power) -> std::io::Result { + let (is_on, watts) = wire::socket_get_status(addr)?; + Ok(Self { + name, + inner: SocketInner::Tcp { + addr, + rated, + state: Mutex::new(TcpState { + is_on, + watts, + last_error: None, + }), + }, + }) + } + + /// Returns `true` if this socket uses TCP to talk to a remote device. + pub fn is_remote(&self) -> bool { + matches!(self.inner, SocketInner::Tcp { .. }) + } + + /// Last I/O or protocol error for TCP mode; `None` for local sockets or after a successful sync. + pub fn last_error(&self) -> Option { + match &self.inner { + SocketInner::Local { .. } => None, + SocketInner::Tcp { state, .. } => state.lock().ok()?.last_error.clone(), + } + } + + fn sync_tcp_cache(addr: SocketAddr, _rated: Power, state: &Mutex) { + let mut guard = match state.lock() { + Ok(g) => g, + Err(_) => return, + }; + match wire::socket_get_status(addr) { + Ok((on, w)) => { + guard.is_on = on; + guard.watts = w; + guard.last_error = None; + } + Err(e) => { + guard.last_error = Some(e.to_string()); + } + } + } + + fn tcp_apply(addr: SocketAddr, state: &Mutex, op: F) + where + F: FnOnce() -> std::io::Result<()>, + { + let mut guard = match state.lock() { + Ok(g) => g, + Err(_) => return, + }; + match op() { + Ok(()) => match wire::socket_get_status(addr) { + Ok((on, w)) => { + guard.is_on = on; + guard.watts = w; + guard.last_error = None; + } + Err(e) => { + guard.last_error = Some(e.to_string()); + } + }, + Err(e) => { + guard.last_error = Some(e.to_string()); + } + } + } + pub fn turn_on(&mut self) { - self.is_on = true; + match &mut self.inner { + SocketInner::Local { is_on, .. } => *is_on = true, + SocketInner::Tcp { addr, state, .. } => { + let addr = *addr; + Self::tcp_apply(addr, state, || wire::socket_set_on(addr)); + } + } } - /// Turn off the socket - /// - /// # Examples - /// - /// ``` - /// use smart_home::{Socket, Power}; - /// - /// let mut socket = Socket::new("TV".to_string(), true, Power::new(120.0).unwrap()); - /// socket.turn_off(); - /// assert!(!socket.is_on()); - /// ``` pub fn turn_off(&mut self) { - self.is_on = false; + match &mut self.inner { + SocketInner::Local { is_on, .. } => *is_on = false, + SocketInner::Tcp { addr, state, .. } => { + let addr = *addr; + Self::tcp_apply(addr, state, || wire::socket_set_off(addr)); + } + } } - /// Check current state (on/off) - /// - /// # Examples - /// - /// ``` - /// use smart_home::{Socket, Power}; - /// - /// let socket = Socket::new("Fridge".to_string(), true, Power::new(200.0).unwrap()); - /// assert!(socket.is_on()); - /// ``` pub fn is_on(&self) -> bool { - self.is_on + match &self.inner { + SocketInner::Local { is_on, .. } => *is_on, + SocketInner::Tcp { addr, rated, state } => { + let addr = *addr; + let rated = *rated; + Self::sync_tcp_cache(addr, rated, state); + state.lock().map(|g| g.is_on).unwrap_or(false) + } + } } - /// Returns current power: zero when off, otherwise the configured power - /// - /// # Examples - /// - /// ``` - /// use smart_home::{Socket, Power}; - /// - /// let socket = Socket::new("Iron".to_string(), true, Power::new(2000.0).unwrap()); - /// assert_eq!(socket.power().watts(), 2000.0); - /// - /// let mut socket_off = Socket::new("Iron".to_string(), false, Power::new(2000.0).unwrap()); - /// assert_eq!(socket_off.power().watts(), 0.0); - /// ``` pub fn power(&self) -> Power { - if self.is_on { - self.power - } else { - Power::zero() + match &self.inner { + SocketInner::Local { is_on, power } => { + if *is_on { + *power + } else { + Power::zero() + } + } + SocketInner::Tcp { addr, rated, state } => { + let addr = *addr; + let rated = *rated; + Self::sync_tcp_cache(addr, rated, state); + let watts = state.lock().map(|g| g.watts).unwrap_or(0.0); + Power::new(watts).unwrap_or_default() + } } } } +impl Clone for Socket { + fn clone(&self) -> Self { + match &self.inner { + SocketInner::Local { is_on, power } => Self { + name: self.name.clone(), + inner: SocketInner::Local { + is_on: *is_on, + power: *power, + }, + }, + SocketInner::Tcp { addr, rated, state } => Self { + name: self.name.clone(), + inner: SocketInner::Tcp { + addr: *addr, + rated: *rated, + state: Mutex::new(state.lock().map(|g| g.clone()).unwrap_or(TcpState { + is_on: false, + watts: 0.0, + last_error: Some("mutex poisoned".to_string()), + })), + }, + }, + } + } +} + +impl PartialEq for Socket { + fn eq(&self, other: &Self) -> bool { + if self.name != other.name { + return false; + } + match (&self.inner, &other.inner) { + ( + SocketInner::Local { + is_on: a, + power: p1, + }, + SocketInner::Local { + is_on: b, + power: p2, + }, + ) => a == b && p1 == p2, + ( + SocketInner::Tcp { + addr: aa, + rated: r1, + .. + }, + SocketInner::Tcp { + addr: ab, + rated: r2, + .. + }, + ) => aa == ab && r1 == r2, + _ => false, + } + } +} + +impl fmt::Display for Socket { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "Socket({})", self.name) + } +} + impl DeviceInfo for Socket { fn name(&self) -> &str { &self.name } fn state(&self) -> String { - format!( + let err = self.last_error(); + let base = format!( "Socket '{}': {} (power: {:.1} W)", self.name, - if self.is_on { "on" } else { "off" }, + if self.is_on() { "on" } else { "off" }, self.power().watts() - ) + ); + match err { + Some(e) => format!("{base} [error: {e}]"), + None => base, + } + } +} + +#[cfg(test)] +mod tests { + use std::io::{Read, Write}; + use std::net::TcpListener; + use std::sync::mpsc; + use std::thread; + use std::time::Duration; + + use super::*; + + #[test] + fn tcp_socket_turn_on_off() { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + listener.set_nonblocking(true).unwrap(); + let addr = listener.local_addr().unwrap(); + let flag = std::sync::Arc::new(std::sync::Mutex::new(false)); + let flag2 = flag.clone(); + let (tx, rx) = mpsc::channel::<()>(); + thread::spawn(move || { + for _ in 0..32 { + match listener.accept() { + Ok((stream, _)) => { + let mut stream = stream; + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + stream + .set_write_timeout(Some(Duration::from_secs(2))) + .unwrap(); + let mut line = String::new(); + let mut buf = [0_u8; 1]; + while let Ok(n) = stream.read(&mut buf) { + if n == 0 { + break; + } + if buf[0] == b'\n' { + let l = line.trim(); + if l == "GET_STATUS" { + let on = *flag2.lock().unwrap(); + let w = if on { 99.0 } else { 0.0 }; + let msg = wire::format_status_line(on, w); + let _ = stream.write_all(msg.as_bytes()); + } else if l == "SET_ON" { + *flag2.lock().unwrap() = true; + let _ = stream.write_all(b"OK\n"); + } else if l == "SET_OFF" { + *flag2.lock().unwrap() = false; + let _ = stream.write_all(b"OK\n"); + } + line.clear(); + } else { + line.push(buf[0] as char); + } + } + } + Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => { + if rx.try_recv().is_ok() { + break; + } + thread::sleep(Duration::from_millis(5)); + } + Err(_) => break, + } + } + }); + + thread::sleep(Duration::from_millis(50)); + + let rated = Power::new(99.0).unwrap(); + let mut socket = Socket::connect_tcp("s".to_string(), addr, rated).expect("connect"); + assert!(!socket.is_on()); + socket.turn_on(); + assert!(socket.is_on()); + assert!((socket.power().watts() - 99.0).abs() < 0.01); + socket.turn_off(); + assert!(!socket.is_on()); + let _ = tx.send(()); } } diff --git a/src/devices/thermometer.rs b/src/devices/thermometer.rs index 9286316..049c1bc 100644 --- a/src/devices/thermometer.rs +++ b/src/devices/thermometer.rs @@ -1,47 +1,196 @@ +use std::fmt; +use std::io; +use std::net::{SocketAddr, UdpSocket}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex}; +use std::thread::JoinHandle; +use std::time::Duration; + use super::device::DeviceInfo; use crate::types::Temperature; +use crate::wire; + +/// Background UDP receiver shared by cloned handles; joins its thread when the last `Arc` is dropped. +struct UdpRecv { + last: Arc>>, + received: Arc, + stop: Arc, + handle: Mutex>>, +} + +impl UdpRecv { + fn spawn(bind: SocketAddr) -> io::Result> { + let socket = UdpSocket::bind(bind)?; + socket.set_read_timeout(Some(Duration::from_millis(200)))?; + let last = Arc::new(Mutex::new(None)); + let received = Arc::new(AtomicBool::new(false)); + let stop = Arc::new(AtomicBool::new(false)); + + let last_t = Arc::clone(&last); + let recv_t = Arc::clone(&received); + let stop_t = Arc::clone(&stop); + let handle = std::thread::spawn(move || { + let mut buf = [0_u8; 256]; + loop { + if stop_t.load(Ordering::SeqCst) { + break; + } + match socket.recv_from(&mut buf) { + Ok((n, _)) => { + if let Some(t) = wire::decode_temperature_celsius(&buf[..n]) { + if let Ok(mut g) = last_t.lock() { + *g = Some(t); + } + recv_t.store(true, Ordering::SeqCst); + } + } + Err(ref e) + if e.kind() == io::ErrorKind::WouldBlock + || e.kind() == io::ErrorKind::TimedOut => + { + continue; + } + Err(_) => continue, + } + } + }); + + Ok(Arc::new(Self { + last, + received, + stop, + handle: Mutex::new(Some(handle)), + })) + } +} + +impl Drop for UdpRecv { + fn drop(&mut self) { + self.stop.store(true, Ordering::SeqCst); + if let Ok(mut g) = self.handle.lock() + && let Some(h) = g.take() + { + let _ = h.join(); + } + } +} -/// Thermometer for measuring temperature -#[derive(Debug, Clone, PartialEq)] +enum ThermometerInner { + Local { + temperature: Temperature, + }, + Udp { + initial: Temperature, + shared: Arc, + }, +} + +impl fmt::Debug for ThermometerInner { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + ThermometerInner::Local { temperature } => f + .debug_struct("Local") + .field("temperature", temperature) + .finish(), + ThermometerInner::Udp { initial, shared } => f + .debug_struct("Udp") + .field("initial", initial) + .field("shared", &format_args!("{:p}", Arc::as_ptr(shared))) + .finish(), + } + } +} + +/// Thermometer: local value or UDP-fed last reading in a background thread. pub struct Thermometer { name: String, - temperature: Temperature, + inner: ThermometerInner, } impl Thermometer { - /// Creates a new thermometer - /// - /// # Arguments - /// - /// * `name` - Thermometer name - /// * `temperature` - Initial temperature - /// - /// # Examples - /// - /// ``` - /// use smart_home::{Thermometer, Temperature, DeviceInfo}; - /// - /// let thermometer = Thermometer::new("Kitchen".to_string(), Temperature::celsius(22.5)); - /// assert_eq!(thermometer.name(), "Kitchen"); - /// assert_eq!(thermometer.temperature().as_celsius(), 22.5); - /// ``` + /// Local thermometer (no network thread). pub fn new(name: String, temperature: Temperature) -> Self { - Self { name, temperature } + Self { + name, + inner: ThermometerInner::Local { temperature }, + } } - /// Returns current temperature - /// - /// # Examples - /// - /// ``` - /// use smart_home::{Thermometer, Temperature}; + /// Listens for UDP temperature packets on `bind` in a background thread. /// - /// let thermometer = Thermometer::new("Living Room".to_string(), Temperature::celsius(24.0)); - /// let temp = thermometer.temperature(); - /// assert_eq!(temp.as_celsius(), 24.0); - /// ``` + /// Until the first packet arrives, [`Self::temperature`] returns `initial`. + pub fn bind_udp(name: String, bind: SocketAddr, initial: Temperature) -> io::Result { + let shared = UdpRecv::spawn(bind)?; + Ok(Self { + name, + inner: ThermometerInner::Udp { initial, shared }, + }) + } + + pub fn is_udp(&self) -> bool { + matches!(self.inner, ThermometerInner::Udp { .. }) + } + + /// `true` if at least one UDP datagram was decoded successfully. + pub fn has_udp_reading(&self) -> bool { + match &self.inner { + ThermometerInner::Local { .. } => true, + ThermometerInner::Udp { shared, .. } => shared.received.load(Ordering::SeqCst), + } + } + pub fn temperature(&self) -> Temperature { - self.temperature + match &self.inner { + ThermometerInner::Local { temperature } => *temperature, + ThermometerInner::Udp { initial, shared } => { + shared.last.lock().ok().and_then(|g| *g).unwrap_or(*initial) + } + } + } +} + +impl Clone for Thermometer { + fn clone(&self) -> Self { + Self { + name: self.name.clone(), + inner: match &self.inner { + ThermometerInner::Local { temperature } => ThermometerInner::Local { + temperature: *temperature, + }, + ThermometerInner::Udp { initial, shared } => ThermometerInner::Udp { + initial: *initial, + shared: Arc::clone(shared), + }, + }, + } + } +} + +impl PartialEq for Thermometer { + fn eq(&self, other: &Self) -> bool { + if self.name != other.name { + return false; + } + match (&self.inner, &other.inner) { + ( + ThermometerInner::Local { temperature: a }, + ThermometerInner::Local { temperature: b }, + ) => a == b, + ( + ThermometerInner::Udp { shared: s1, .. }, + ThermometerInner::Udp { shared: s2, .. }, + ) => Arc::ptr_eq(s1, s2), + _ => false, + } + } +} + +impl fmt::Debug for Thermometer { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("Thermometer") + .field("name", &self.name) + .field("inner", &self.inner) + .finish() } } @@ -54,7 +203,41 @@ impl DeviceInfo for Thermometer { format!( "Thermometer '{}': {:.1}°C", self.name, - self.temperature.as_celsius() + self.temperature().as_celsius() ) } } + +#[cfg(test)] +mod tests { + use std::net::UdpSocket; + use std::thread; + use std::time::Duration; + + use super::*; + + #[test] + fn udp_thermometer_receives() { + let addr: SocketAddr = { + let sock = UdpSocket::bind("127.0.0.1:0").unwrap(); + sock.local_addr().unwrap() + }; + + let t = Thermometer::bind_udp("t".to_string(), addr, Temperature::celsius(1.0)).unwrap(); + assert!(!t.has_udp_reading()); + assert!((t.temperature().as_celsius() - 1.0).abs() < 0.01); + + let sender = UdpSocket::bind("127.0.0.1:0").unwrap(); + let payload = wire::encode_temperature_celsius(19.25); + sender.send_to(&payload, addr).unwrap(); + + for _ in 0..50 { + if t.has_udp_reading() { + break; + } + thread::sleep(Duration::from_millis(10)); + } + assert!(t.has_udp_reading()); + assert!((t.temperature().as_celsius() - 19.25).abs() < 0.01); + } +} diff --git a/src/lib.rs b/src/lib.rs index 97a45e9..972bb59 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -7,6 +7,7 @@ //! - [`types`] — [`Power`], [`Temperature`] //! - [`report`] — [`Report`] //! - [`error`] — [`SmartHomeError`] +//! - [`wire`] — TCP/UDP framing helpers shared with `socket_sim` / `thermo_sim` //! //! Types are re-exported at the crate root. //! @@ -33,6 +34,7 @@ pub mod error; pub mod home; pub mod report; pub mod types; +pub mod wire; pub use devices::{Device, DeviceInfo, Socket, Thermometer}; pub use error::SmartHomeError; diff --git a/src/wire.rs b/src/wire.rs new file mode 100644 index 0000000..4eecc6a --- /dev/null +++ b/src/wire.rs @@ -0,0 +1,151 @@ +//! Line-oriented TCP commands for smart sockets and UDP payload for thermometers. + +use crate::types::Temperature; +use std::io::{self, Read, Write}; +use std::net::{SocketAddr, TcpStream, ToSocketAddrs}; +use std::time::Duration; + +/// Default timeout for synchronous TCP operations from library code. +pub const TCP_TIMEOUT: Duration = Duration::from_secs(2); + +pub fn tcp_connect(addr: impl ToSocketAddrs) -> io::Result { + let stream = TcpStream::connect_timeout( + &addr + .to_socket_addrs()? + .next() + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "empty address list"))?, + TCP_TIMEOUT, + )?; + stream.set_read_timeout(Some(TCP_TIMEOUT))?; + stream.set_write_timeout(Some(TCP_TIMEOUT))?; + Ok(stream) +} + +pub fn read_line(stream: &mut TcpStream) -> io::Result { + let mut buf = Vec::new(); + let mut byte = [0_u8; 1]; + loop { + let n = stream.read(&mut byte)?; + if n == 0 { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "connection closed before newline", + )); + } + if byte[0] == b'\n' { + break; + } + buf.push(byte[0]); + } + String::from_utf8(buf).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e)) +} + +pub fn write_line(stream: &mut TcpStream, line: &[u8]) -> io::Result<()> { + stream.write_all(line)?; + if !line.ends_with(b"\n") { + stream.write_all(b"\n")?; + } + stream.flush() +} + +pub fn socket_get_status(addr: SocketAddr) -> io::Result<(bool, f32)> { + let mut stream = tcp_connect(addr)?; + write_line(&mut stream, b"GET_STATUS")?; + let line = read_line(&mut stream)?; + parse_status_line(&line) +} + +pub fn socket_set_on(addr: SocketAddr) -> io::Result<()> { + let mut stream = tcp_connect(addr)?; + write_line(&mut stream, b"SET_ON")?; + expect_ok(&mut stream) +} + +pub fn socket_set_off(addr: SocketAddr) -> io::Result<()> { + let mut stream = tcp_connect(addr)?; + write_line(&mut stream, b"SET_OFF")?; + expect_ok(&mut stream) +} + +fn expect_ok(stream: &mut TcpStream) -> io::Result<()> { + let line = read_line(stream)?; + if line.trim() == "OK" { + Ok(()) + } else { + Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("expected OK, got {line:?}"), + )) + } +} + +fn parse_status_line(line: &str) -> io::Result<(bool, f32)> { + let line = line.trim(); + let rest = line + .strip_prefix("STATUS,") + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "bad STATUS prefix"))?; + let mut parts = rest.splitn(2, ','); + let flag = parts + .next() + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "STATUS missing on/off field"))?; + let watts = parts + .next() + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "STATUS missing watts"))? + .parse::() + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + let on = match flag { + "1" | "on" | "ON" | "true" | "TRUE" => true, + "0" | "off" | "OFF" | "false" | "FALSE" => false, + _ => { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("bad on/off flag: {flag}"), + )); + } + }; + Ok((on, watts)) +} + +pub fn format_status_line(is_on: bool, watts: f32) -> String { + let flag = if is_on { "1" } else { "0" }; + format!("STATUS,{flag},{watts}\n") +} + +/// UDP payload: 4-byte little-endian IEEE754 `f32` (degrees Celsius). +pub fn encode_temperature_celsius(c: f32) -> [u8; 4] { + c.to_le_bytes() +} + +pub fn decode_temperature_celsius(buf: &[u8]) -> Option { + if buf.len() < 4 { + return None; + } + let mut arr = [0_u8; 4]; + arr.copy_from_slice(&buf[..4]); + let v = f32::from_le_bytes(arr); + Some(Temperature::celsius(v)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parse_status_roundtrip() { + let s = "STATUS,1,123.5"; + let (on, w) = parse_status_line(s).unwrap(); + assert!(on); + assert!((w - 123.5).abs() < 0.01); + let line = format_status_line(true, 123.5); + let (on2, w2) = parse_status_line(line.trim_end()).unwrap(); + assert!(on2); + assert!((w2 - 123.5).abs() < 0.01); + } + + #[test] + fn temperature_udp_roundtrip() { + let b = encode_temperature_celsius(-3.5); + let t = decode_temperature_celsius(&b).unwrap(); + assert!((t.as_celsius() + 3.5).abs() < 0.0001); + } +}