From cf50df31c5d971fc2e574c4bd7000be7e0f79a73 Mon Sep 17 00:00:00 2001 From: Janis Hutz Date: Mon, 24 Aug 2026 17:33:09 +0200 Subject: [PATCH] feat: working mqtt integration, basic GPIO initialization --- config.yml | 1 + src/conf.rs | 10 ++++- src/gpio_utils.rs | 99 ++++++++++++++++++++++++++++++++++++++--------- src/main.rs | 3 +- src/mqtt.rs | 35 +++++++++++------ 5 files changed, 115 insertions(+), 33 deletions(-) diff --git a/config.yml b/config.yml index 1af8811..48ab05d 100644 --- a/config.yml +++ b/config.yml @@ -12,5 +12,6 @@ topics: - topic: "" pin: 0 mode: "in" + offTimeout: -1 # -1 disables automatic switch off of the pin pollInterval: 100 diff --git a/src/conf.rs b/src/conf.rs index 2066c5b..1a92861 100644 --- a/src/conf.rs +++ b/src/conf.rs @@ -6,6 +6,7 @@ pub struct Topic { pub topic: String, pub pin: u16, pub mode: PinMode, + pub off_timeout: u64, } #[derive(Debug, PartialEq, Clone, Copy)] @@ -27,7 +28,7 @@ pub struct MqttConfig { pub struct Config { pub mqtt: MqttConfig, pub topics: Vec, - pub poll_interval: u16, + pub poll_interval: u64, } pub fn load_config(path: &str) -> Result> { @@ -69,6 +70,11 @@ pub fn load_config(path: &str) -> Result> { .expect(&format!("Invalid topic configuration found at index {}. A pin is required", i)) .as_i64() .expect(&format!("Invalid topic configuration found at index {}. The pin should be a 16 bit integer (0-65535)", i)) as u16, + off_timeout: topic + .get(&Yaml::String(String::from("pin"))) + .unwrap_or(&Yaml::Integer(-1)) + .as_i64() + .expect(&format!("Invalid topic configuration found at index {}. The pin should be a 16 bit integer (0-65535)", i)) as u64, mode: if topic .get(&Yaml::String(String::from("mode"))) .expect(&format!("Invalid topic configuration found at index {}. Mode is unset", i)) @@ -119,6 +125,6 @@ pub fn load_config(path: &str) -> Result> { .get(&Yaml::String(String::from("pollInterval"))) .expect("Missing config for pollInterval") .as_i64() - .expect("pollInterval is not an integer value") as u16, + .expect("pollInterval is not an integer value") as u64, }) } diff --git a/src/gpio_utils.rs b/src/gpio_utils.rs index afcaac6..e779b84 100644 --- a/src/gpio_utils.rs +++ b/src/gpio_utils.rs @@ -1,52 +1,115 @@ +use gpio::{GpioIn, GpioOut}; +use rumqttc::Publish; use std::collections::HashMap; use crate::conf::{PinMode, Topic}; +pub struct InputPin { + pub gpio: gpio::sysfs::SysFsGpioInput, + pub id: u16, + pub topic: String, +} + +pub struct OutputPin { + gpio: gpio::sysfs::SysFsGpioOutput, + off_timeout: u64, + pub id: u16, + pub topic: String, +} + /// Split topics into in and out pins. Returns them in this order /// /// * `topics`: The topics to use -pub fn split_topics_and_configure_pins(topics: Vec) -> (Vec, Vec) { - let mut in_pins: Vec = Vec::new(); - let mut out_pins: Vec = Vec::new(); +pub fn split_topics_and_configure_pins(topics: Vec) -> (Vec, Vec) { + let mut in_pins: Vec = Vec::new(); + let mut out_pins: Vec = Vec::new(); + let mut in_used: Vec = Vec::new(); + let mut out_used: Vec = Vec::new(); for topic in topics { if topic.mode == PinMode::IN { - in_pins.push(topic); + if in_used.contains(&topic.pin) { + println!( + "Warning: Pin {} used more than once, all further uses dropped", + topic.pin + ); + continue; + } else if out_used.contains(&topic.pin) { + panic!("Pin {} used for both input and output!", topic.pin) + } + in_pins.push(InputPin { + topic: topic.topic, + gpio: gpio::sysfs::SysFsGpioInput::open(topic.pin) + .expect(&format!("Unable to find GPIO pin {}", topic.pin)), + id: topic.pin, + }); + in_used.push(topic.pin) } else { - out_pins.push(topic); + if out_used.contains(&topic.pin) { + println!( + "Warning: Pin {} used more than once, all further uses dropped", + topic.pin + ); + continue; + } else if in_used.contains(&topic.pin) { + panic!("Pin {} used for both input and output!", topic.pin) + } + out_pins.push(OutputPin { + topic: topic.topic, + gpio: gpio::sysfs::SysFsGpioOutput::open(topic.pin) + .expect(&format!("Unable to find GPIO pin {}", topic.pin)), + id: topic.pin, + off_timeout: topic.off_timeout, + }); + out_used.push(topic.pin) } } return (in_pins, out_pins); } -pub fn read_pin_value(id: u16) {} +pub fn read_pin(pin: &mut gpio::sysfs::SysFsGpioInput) -> u8 { + if let gpio::GpioValue::Low = pin.read_value().expect("Failed reading value") { + return 1; + } else { + return 0; + } +} pub struct GPIOController { - in_pins: HashMap, + out_pins: HashMap, + topics: Vec, } impl GPIOController { /// Create a new GPIO controller. /// Note that the move of ownership is intentional behaviour. Only create one instance of this controller! /// - /// * `topics`: The topics that are going to be managed by this controller (All OUT mode pins - /// are dropped without error) - pub fn new(topics: Vec) -> Self { + /// * `pins`: The pins that are going to be managed by this controller + pub fn new(pins: Vec) -> Self { let mut controller = GPIOController { - in_pins: HashMap::new(), + out_pins: HashMap::new(), + topics: Vec::new(), }; // Build hash maps - for topic in topics { - if topic.mode == PinMode::IN { - controller - .in_pins - .insert(String::from(topic.topic.as_str()), topic); - } + for pin in pins { + let topic = pin.topic.clone(); + controller.out_pins.insert(pin.topic.clone(), pin); + controller.topics.push(topic); } return controller; } - pub fn handle_event(&self) {} + pub fn handle_event(&mut self, instruction: &Publish) { + if self.topics.contains(&instruction.topic) { + let pin = self + .out_pins + .get(&instruction.topic) + .expect("Failed to load pins data"); + println!("Hello World, pin {}", pin.id); + println!("Payload {}", instruction.payload.first().unwrap()); + // TODO: Off timeout + } + } } diff --git a/src/main.rs b/src/main.rs index fc62dad..fd79686 100644 --- a/src/main.rs +++ b/src/main.rs @@ -8,7 +8,8 @@ fn main() { // CLI interface should allow setting the config file println!("mqtt-remote-gpio"); - let config = conf::load_config("config.yml").unwrap(); + // let config = conf::load_config("config.yml").unwrap(); + let config = conf::load_config("config.secret.yml").unwrap(); // TODO: Remove this when done println!("{:#?}", config); diff --git a/src/mqtt.rs b/src/mqtt.rs index 5ba4392..6b79c8e 100644 --- a/src/mqtt.rs +++ b/src/mqtt.rs @@ -1,6 +1,7 @@ -use crate::conf::{Config}; -use crate::gpio_utils::{GPIOController, split_topics_and_configure_pins}; -use rumqttc::{Client, MqttOptions, QoS}; +use crate::conf::Config; +use crate::gpio_utils::{GPIOController, read_pin, split_topics_and_configure_pins}; +use rumqttc::Packet::Publish; +use rumqttc::{Client, Event, MqttOptions, QoS}; use std::{thread, time::Duration}; /// MQTT connection handler @@ -9,7 +10,7 @@ use std::{thread, time::Duration}; /// * `mqttoptions`: MQTT options for conenction pub fn handler(config: Config) { // Configure pins - let (subscribe_topics, publish_topics) = split_topics_and_configure_pins(config.topics); + let (mut input_pins, output_pins) = split_topics_and_configure_pins(config.topics); // Configure MQTT let mut mqttoptions = MqttOptions::new("mqtt", config.mqtt.host, config.mqtt.port); @@ -20,7 +21,7 @@ pub fn handler(config: Config) { let (client, mut connection) = Client::new(mqttoptions, 10); // Subscribe to events - for topic in subscribe_topics.iter() { + for topic in output_pins.iter() { println!("Setting up topic {}", topic.topic); client .subscribe(topic.topic.as_str(), QoS::AtMostOnce) @@ -30,22 +31,32 @@ pub fn handler(config: Config) { // Spawn thread to monitor pins thread::spawn(move || { loop { - for topic in &publish_topics { + for topic in input_pins.iter_mut() { // TODO: handle fails client - .publish(&topic.topic, QoS::AtLeastOnce, false, vec![0 as u8; 1]) + .publish( + &topic.topic, + QoS::AtLeastOnce, + false, + vec![read_pin(&mut topic.gpio); 1], + ) .unwrap_or_default(); } - thread::sleep(Duration::from_millis(100)); + thread::sleep(Duration::from_millis(config.poll_interval)); } }); // Main thread listens to topic updates - let controller = GPIOController::new(subscribe_topics); + let mut controller = GPIOController::new(output_pins); for (_, notification) in connection.iter().enumerate() { - println!("Notification {:?}", notification); - controller.handle_event(); - thread::sleep(Duration::from_millis(100)); + let msg = notification.unwrap(); + if let Event::Incoming(val) = msg { + if let Publish(content) = val { + controller.handle_event(&content); + } + } + // println!("Notification {:?}", msg); + thread::sleep(Duration::from_millis(config.poll_interval / 2)); // TODO: Handle pin value change instructions } }