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
3 changes: 1 addition & 2 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,3 +1,2 @@
/target
/tmp
/examples
/tmp
182 changes: 182 additions & 0 deletions examples/simulated_home.rs
Original file line number Diff line number Diff line change
@@ -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<String> = 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<String, Device> = 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<String, Device> = 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}");
}
}
}
3 changes: 3 additions & 0 deletions examples/thermo_kitchen.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
127.0.0.1:19101
300
21.0
3 changes: 3 additions & 0 deletions examples/thermo_living.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
127.0.0.1:19102
300
23.5
187 changes: 187 additions & 0 deletions src/bin/socket_sim.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,187 @@
//! TCP smart-outlet simulator: non-blocking accept/read, shared state, many clients.
//!
//! Usage: `socket_sim <BIND_ADDR> [--watts <W>] [--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<u8>,
buf: Vec<u8>,
}

impl Client {
fn new(stream: TcpStream) -> io::Result<Self> {
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<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 {
"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<dyn std::error::Error>> {
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<Client> = 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));
}
}
Loading