Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 76 additions & 0 deletions examples/patterns.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
use std::cell::RefCell;
use std::rc::Rc;

use smart_home::{
Device, DeviceInfo, HomeBuilder, Reporter, Room, Socket, Subscriber, Temperature, Thermometer,
};

#[derive(Default)]
struct MySubscriber {
added: Rc<RefCell<Vec<String>>>,
}

impl MySubscriber {
fn with_log(added: Rc<RefCell<Vec<String>>>) -> Self {
Self { added }
}
}

impl Subscriber for MySubscriber {
fn on_event(&mut self, device: &Device) {
self.added.borrow_mut().push(device.name().to_string());
}
}

fn main() {
let home = HomeBuilder::new()
.add_room("First room")
.add_device("Socket_1", Socket::default())
.add_device("Socket_2", Socket::default())
.add_device("Thermo_1", Thermometer::default())
.add_room("Second room")
.add_device("Socket_3", Socket::default())
.add_device(
"Thermo_2",
Thermometer::new("Thermo_2".to_string(), Temperature::celsius(24.0)),
)
.build();

println!("=== HomeBuilder report ===");
Reporter::new().add(&home).report();

let mut room = Room::default();
let subscriber_log = Rc::new(RefCell::new(Vec::new()));
room.subscribe(MySubscriber::with_log(Rc::clone(&subscriber_log)));

let closure_log = Rc::new(RefCell::new(Vec::new()));
let closure_log_handle = Rc::clone(&closure_log);
room.subscribe(move |device: &Device| {
closure_log_handle
.borrow_mut()
.push(format!("closure: {}", device.name()));
});

room.insert_device("Socket_4".to_string(), Socket::default().into());
room.insert_device("Thermo_3".to_string(), Thermometer::default().into());

println!("\n=== Observer logs ===");
println!("subscriber: {:?}", subscriber_log.borrow());
println!("closure: {:?}", closure_log.borrow());

let device = Device::default();
let socket1 = Socket::default();
let socket2 = Socket::default();
let thermo1 = Thermometer::default();
let thermo2 = Thermometer::new("Thermo".to_string(), Temperature::celsius(21.5));

println!("\n=== Reporter composite ===");
Reporter::new()
.add(&room)
.add(&device)
.add(&socket1)
.add(&socket2)
.add(&thermo1)
.add(&thermo2)
.report();
}
42 changes: 28 additions & 14 deletions src/bin/socket_sim.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
//! TCP smart-outlet simulator: non-blocking accept/read, shared state, many clients.
//! TCP smart-outlet simulator: non-blocking accept/read, length-prefixed frames, shared state, many clients.
//!
//! Usage: `socket_sim <BIND_ADDR> [--watts <W>] [--on|--off]`

Expand Down Expand Up @@ -37,11 +37,11 @@ impl Client {
})
}

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 push_response(&mut self, payload: &str) {
let bytes = payload.as_bytes();
let len = u32::try_from(bytes.len()).expect("socket_sim response frame is too large");
self.out.extend_from_slice(&len.to_be_bytes());
self.out.extend_from_slice(bytes);
}

fn flush_writes(&mut self) -> io::Result<()> {
Expand Down Expand Up @@ -78,18 +78,32 @@ impl Client {
Ok(())
}

fn drain_lines(&mut self, state: &Arc<Mutex<OutletState>>) -> io::Result<()> {
while let Some(pos) = self.buf.iter().position(|&b| b == b'\n') {
let raw: Vec<u8> = self.buf.drain(..=pos).collect();
let line = String::from_utf8_lossy(&raw[..raw.len().saturating_sub(1)]);
let line = line.trim();
match line {
fn drain_frames(&mut self, state: &Arc<Mutex<OutletState>>) -> io::Result<()> {
while self.buf.len() >= 4 {
let len =
u32::from_be_bytes(self.buf[..4].try_into().expect("length prefix is 4 bytes"))
as usize;
if len > 1024 {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("frame is too large: {len} bytes"),
));
}
if self.buf.len() < 4 + len {
break;
}
self.buf.drain(..4);
let raw: Vec<u8> = self.buf.drain(..len).collect();
let command = String::from_utf8(raw).map_err(|e| {
io::Error::new(io::ErrorKind::InvalidData, format!("bad UTF-8: {e}"))
})?;
match command.as_str() {
"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());
self.push_response(&msg);
}
"SET_ON" => {
if let Ok(mut g) = state.lock() {
Expand Down Expand Up @@ -172,7 +186,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
let c = &mut clients[i];
match c.read_available() {
Err(_) => true,
Ok(()) => c.drain_lines(&state).is_err() || c.flush_writes().is_err(),
Ok(()) => c.drain_frames(&state).is_err() || c.flush_writes().is_err(),
}
};
if remove {
Expand Down
6 changes: 6 additions & 0 deletions src/devices/device.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,12 @@ impl From<Socket> for Device {
}
}

impl Default for Device {
fn default() -> Self {
Device::Socket(Socket::default())
}
}

impl Report for Device {
fn report(&self) -> String {
format!("{}\n", DeviceInfo::state(self))
Expand Down
49 changes: 25 additions & 24 deletions src/devices/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ use std::net::SocketAddr;
use std::sync::Mutex;

use super::device::DeviceInfo;
use crate::report::Report;
use crate::types::Power;
use crate::wire;

Expand Down Expand Up @@ -195,6 +196,12 @@ impl Clone for Socket {
}
}

impl Default for Socket {
fn default() -> Self {
Self::new("Socket".to_string(), false, Power::default())
}
}

impl PartialEq for Socket {
fn eq(&self, other: &Self) -> bool {
if self.name != other.name {
Expand Down Expand Up @@ -254,9 +261,14 @@ impl DeviceInfo for Socket {
}
}

impl Report for Socket {
fn report(&self) -> String {
format!("{}\n", self.state())
}
}

#[cfg(test)]
mod tests {
use std::io::{Read, Write};
use std::net::TcpListener;
use std::sync::mpsc;
use std::thread;
Expand All @@ -283,29 +295,18 @@ mod tests {
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);
while let Ok(command) = wire::read_frame(&mut stream) {
if command == "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 _ = wire::write_frame(&mut stream, msg.as_bytes());
} else if command == "SET_ON" {
*flag2.lock().unwrap() = true;
let _ = wire::write_frame(&mut stream, b"OK");
} else if command == "SET_OFF" {
*flag2.lock().unwrap() = false;
let _ = wire::write_frame(&mut stream, b"OK");
}
}
}
Expand Down
13 changes: 13 additions & 0 deletions src/devices/thermometer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ use std::thread::JoinHandle;
use std::time::Duration;

use super::device::DeviceInfo;
use crate::report::Report;
use crate::types::Temperature;
use crate::wire;

Expand Down Expand Up @@ -166,6 +167,12 @@ impl Clone for Thermometer {
}
}

impl Default for Thermometer {
fn default() -> Self {
Self::new("Thermometer".to_string(), Temperature::default())
}
}

impl PartialEq for Thermometer {
fn eq(&self, other: &Self) -> bool {
if self.name != other.name {
Expand Down Expand Up @@ -208,6 +215,12 @@ impl DeviceInfo for Thermometer {
}
}

impl Report for Thermometer {
fn report(&self) -> String {
format!("{}\n", self.state())
}
}

#[cfg(test)]
mod tests {
use std::net::UdpSocket;
Expand Down
88 changes: 88 additions & 0 deletions src/home/builder.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
use std::collections::HashMap;
use std::marker::PhantomData;

use crate::devices::Device;

use super::{Room, SmartHome};

/// Builder state before the first room is added.
pub struct NoRoom;

/// Builder state after at least one room is available.
pub struct WithRoom;

/// Typestate builder for [`SmartHome`].
pub struct HomeBuilder<State = NoRoom> {
name: String,
rooms: HashMap<String, Room>,
current_room: Option<String>,
_state: PhantomData<State>,
}

impl HomeBuilder<NoRoom> {
pub fn new() -> Self {
Self {
name: "Smart Home".to_string(),
rooms: HashMap::new(),
current_room: None,
_state: PhantomData,
}
}

pub fn with_name(name: impl Into<String>) -> Self {
Self {
name: name.into(),
rooms: HashMap::new(),
current_room: None,
_state: PhantomData,
}
}

pub fn add_room(mut self, name: impl Into<String>) -> HomeBuilder<WithRoom> {
let name = name.into();
self.rooms
.insert(name.clone(), Room::new(name.clone(), HashMap::new()));
HomeBuilder {
name: self.name,
rooms: self.rooms,
current_room: Some(name),
_state: PhantomData,
}
}

pub fn build(self) -> SmartHome {
SmartHome::new(self.name, self.rooms)
}
}

impl Default for HomeBuilder<NoRoom> {
fn default() -> Self {
Self::new()
}
}

impl HomeBuilder<WithRoom> {
pub fn add_room(mut self, name: impl Into<String>) -> Self {
let name = name.into();
self.rooms
.insert(name.clone(), Room::new(name.clone(), HashMap::new()));
self.current_room = Some(name);
self
}

pub fn add_device<D>(mut self, key: impl Into<String>, device: D) -> Self
where
D: Into<Device>,
{
if let Some(current_room) = self.current_room.as_deref()
&& let Some(room) = self.rooms.get_mut(current_room)
{
room.insert_device(key.into(), device.into());
}
self
}

pub fn build(self) -> SmartHome {
SmartHome::new(self.name, self.rooms)
}
}
4 changes: 3 additions & 1 deletion src/home/mod.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
pub mod builder;
pub mod room;
pub mod smart_home;

pub use room::Room;
pub use builder::{HomeBuilder, NoRoom, WithRoom};
pub use room::{Room, Subscriber};
pub use smart_home::SmartHome;
Loading