feat: working mqtt integration, basic GPIO initialization

This commit is contained in:
janishutz committed 2026-08-24 17:33:09 +02:00
1 parent f044a5b5bd
commit cf50df31c5
5 files changed
+115 -33

No files matched your search

+1
View File
@@ -12,5 +12,6 @@ topics:
- topic: "" - topic: ""
pin: 0 pin: 0
mode: "in" mode: "in"
offTimeout: -1 # -1 disables automatic switch off of the pin
pollInterval: 100 pollInterval: 100
+8 -2
View File
@@ -6,6 +6,7 @@ pub struct Topic {
pub topic: String, pub topic: String,
pub pin: u16, pub pin: u16,
pub mode: PinMode, pub mode: PinMode,
pub off_timeout: u64,
} }
#[derive(Debug, PartialEq, Clone, Copy)] #[derive(Debug, PartialEq, Clone, Copy)]
@@ -27,7 +28,7 @@ pub struct MqttConfig {
pub struct Config { pub struct Config {
pub mqtt: MqttConfig, pub mqtt: MqttConfig,
pub topics: Vec<Topic>, pub topics: Vec<Topic>,
pub poll_interval: u16, pub poll_interval: u64,
} }
pub fn load_config(path: &str) -> Result<Config, Box<dyn std::error::Error>> { pub fn load_config(path: &str) -> Result<Config, Box<dyn std::error::Error>> {
@@ -69,6 +70,11 @@ pub fn load_config(path: &str) -> Result<Config, Box<dyn std::error::Error>> {
.expect(&format!("Invalid topic configuration found at index {}. A pin is required", i)) .expect(&format!("Invalid topic configuration found at index {}. A pin is required", i))
.as_i64() .as_i64()
.expect(&format!("Invalid topic configuration found at index {}. The pin should be a 16 bit integer (0-65535)", i)) as u16, .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 mode: if topic
.get(&Yaml::String(String::from("mode"))) .get(&Yaml::String(String::from("mode")))
.expect(&format!("Invalid topic configuration found at index {}. Mode is unset", i)) .expect(&format!("Invalid topic configuration found at index {}. Mode is unset", i))
@@ -119,6 +125,6 @@ pub fn load_config(path: &str) -> Result<Config, Box<dyn std::error::Error>> {
.get(&Yaml::String(String::from("pollInterval"))) .get(&Yaml::String(String::from("pollInterval")))
.expect("Missing config for pollInterval") .expect("Missing config for pollInterval")
.as_i64() .as_i64()
.expect("pollInterval is not an integer value") as u16, .expect("pollInterval is not an integer value") as u64,
}) })
} }
+81 -18
View File
@@ -1,52 +1,115 @@
use gpio::{GpioIn, GpioOut};
use rumqttc::Publish;
use std::collections::HashMap; use std::collections::HashMap;
use crate::conf::{PinMode, Topic}; 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 /// Split topics into in and out pins. Returns them in this order
/// ///
/// * `topics`: The topics to use /// * `topics`: The topics to use
pub fn split_topics_and_configure_pins(topics: Vec<Topic>) -> (Vec<Topic>, Vec<Topic>) { pub fn split_topics_and_configure_pins(topics: Vec<Topic>) -> (Vec<InputPin>, Vec<OutputPin>) {
let mut in_pins: Vec<Topic> = Vec::new(); let mut in_pins: Vec<InputPin> = Vec::new();
let mut out_pins: Vec<Topic> = Vec::new(); let mut out_pins: Vec<OutputPin> = Vec::new();
let mut in_used: Vec<u16> = Vec::new();
let mut out_used: Vec<u16> = Vec::new();
for topic in topics { for topic in topics {
if topic.mode == PinMode::IN { 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 { } 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); 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 { pub struct GPIOController {
in_pins: HashMap<String, Topic>, out_pins: HashMap<String, OutputPin>,
topics: Vec<String>,
} }
impl GPIOController { impl GPIOController {
/// Create a new GPIO controller. /// Create a new GPIO controller.
/// Note that the move of ownership is intentional behaviour. Only create one instance of this 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 /// * `pins`: The pins that are going to be managed by this controller
/// are dropped without error) pub fn new(pins: Vec<OutputPin>) -> Self {
pub fn new(topics: Vec<Topic>) -> Self {
let mut controller = GPIOController { let mut controller = GPIOController {
in_pins: HashMap::new(), out_pins: HashMap::new(),
topics: Vec::new(),
}; };
// Build hash maps // Build hash maps
for topic in topics { for pin in pins {
if topic.mode == PinMode::IN { let topic = pin.topic.clone();
controller controller.out_pins.insert(pin.topic.clone(), pin);
.in_pins controller.topics.push(topic);
.insert(String::from(topic.topic.as_str()), topic);
}
} }
return controller; 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
}
}
} }
+2 -1
View File
@@ -8,7 +8,8 @@ fn main() {
// CLI interface should allow setting the config file // CLI interface should allow setting the config file
println!("mqtt-remote-gpio"); 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 // TODO: Remove this when done
println!("{:#?}", config); println!("{:#?}", config);
+23 -12
View File
@@ -1,6 +1,7 @@
use crate::conf::{Config}; use crate::conf::Config;
use crate::gpio_utils::{GPIOController, split_topics_and_configure_pins}; use crate::gpio_utils::{GPIOController, read_pin, split_topics_and_configure_pins};
use rumqttc::{Client, MqttOptions, QoS}; use rumqttc::Packet::Publish;
use rumqttc::{Client, Event, MqttOptions, QoS};
use std::{thread, time::Duration}; use std::{thread, time::Duration};
/// MQTT connection handler /// MQTT connection handler
@@ -9,7 +10,7 @@ use std::{thread, time::Duration};
/// * `mqttoptions`: MQTT options for conenction /// * `mqttoptions`: MQTT options for conenction
pub fn handler(config: Config) { pub fn handler(config: Config) {
// Configure pins // 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 // Configure MQTT
let mut mqttoptions = MqttOptions::new("mqtt", config.mqtt.host, config.mqtt.port); 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); let (client, mut connection) = Client::new(mqttoptions, 10);
// Subscribe to events // Subscribe to events
for topic in subscribe_topics.iter() { for topic in output_pins.iter() {
println!("Setting up topic {}", topic.topic); println!("Setting up topic {}", topic.topic);
client client
.subscribe(topic.topic.as_str(), QoS::AtMostOnce) .subscribe(topic.topic.as_str(), QoS::AtMostOnce)
@@ -30,22 +31,32 @@ pub fn handler(config: Config) {
// Spawn thread to monitor pins // Spawn thread to monitor pins
thread::spawn(move || { thread::spawn(move || {
loop { loop {
for topic in &publish_topics { for topic in input_pins.iter_mut() {
// TODO: handle fails // TODO: handle fails
client 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(); .unwrap_or_default();
} }
thread::sleep(Duration::from_millis(100)); thread::sleep(Duration::from_millis(config.poll_interval));
} }
}); });
// Main thread listens to topic updates // 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() { for (_, notification) in connection.iter().enumerate() {
println!("Notification {:?}", notification); let msg = notification.unwrap();
controller.handle_event(); if let Event::Incoming(val) = msg {
thread::sleep(Duration::from_millis(100)); 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 // TODO: Handle pin value change instructions
} }
} }