mirror of
https://github.com/janishutz/mqtt-remote-gpio.git
synced 2026-10-09 05:06:21 +02:00
feat: separate files, basic gpio util prep, config parsing
This commit is contained in:
1 parent
f665027373
commit
f044a5b5bd
9 files changed
+290
-167
No files matched your search
+124
@@ -0,0 +1,124 @@
|
||||
use std::fs;
|
||||
use yaml_rust2::{Yaml, YamlLoader};
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct Topic {
|
||||
pub topic: String,
|
||||
pub pin: u16,
|
||||
pub mode: PinMode,
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Clone, Copy)]
|
||||
pub enum PinMode {
|
||||
IN,
|
||||
OUT,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct MqttConfig {
|
||||
pub host: String,
|
||||
pub port: u16,
|
||||
pub authentication: bool,
|
||||
pub user: String,
|
||||
pub password: String,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct Config {
|
||||
pub mqtt: MqttConfig,
|
||||
pub topics: Vec<Topic>,
|
||||
pub poll_interval: u16,
|
||||
}
|
||||
|
||||
pub fn load_config(path: &str) -> Result<Config, Box<dyn std::error::Error>> {
|
||||
// Read config file
|
||||
let conf = fs::read_to_string(path).unwrap_or_else(|_| String::from("test"));
|
||||
let yaml = &YamlLoader::load_from_str(&conf).unwrap()[0];
|
||||
let conf = yaml
|
||||
.as_hash()
|
||||
.expect("Invalid config found at line 1. Not an object");
|
||||
let mqtt_conf = conf
|
||||
.get(&Yaml::String(String::from("mqtt")))
|
||||
.expect("MQTT config is missing.")
|
||||
.as_hash()
|
||||
.expect("MQTT config is invalid");
|
||||
|
||||
// Load topics and pin config
|
||||
let mut topics: Vec<Topic> = Vec::new();
|
||||
let mut i = 0;
|
||||
for raw_topic in conf
|
||||
.get(&Yaml::String(String::from("topics")))
|
||||
.expect("Topics config missing")
|
||||
.as_vec()
|
||||
.expect("Topics config invalid. Expected an array")
|
||||
{
|
||||
i += 1;
|
||||
let topic = raw_topic.as_hash().expect(
|
||||
&format!("Invalid topic configuration found for topic at index {}. All topics should be objects with topic, pin and mode!", i));
|
||||
topics.push(
|
||||
Topic {
|
||||
topic: String::from(
|
||||
topic
|
||||
.get(&Yaml::String(String::from("topic")))
|
||||
.expect(&format!("Invalid topic configuration found at index {}. The topic name is missing", i))
|
||||
.as_str()
|
||||
.expect(&format!("Invalid topic configuration found index {}. The topic name should be a string", i)),
|
||||
),
|
||||
pin: topic
|
||||
.get(&Yaml::String(String::from("pin")))
|
||||
.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,
|
||||
mode: if topic
|
||||
.get(&Yaml::String(String::from("mode")))
|
||||
.expect(&format!("Invalid topic configuration found at index {}. Mode is unset", i))
|
||||
.as_str()
|
||||
.expect(&format!("Invalid topic configuration found at index {}. Mode should be a string of either 'in' or 'out'", i)) == "in"
|
||||
{ PinMode::IN } else { PinMode::OUT },
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
// Create the config struct
|
||||
Ok(Config {
|
||||
mqtt: MqttConfig {
|
||||
host: String::from(
|
||||
mqtt_conf
|
||||
.get(&Yaml::String(String::from("host")))
|
||||
.expect("Host config missing")
|
||||
.as_str()
|
||||
.expect("Host configuration is invalid"),
|
||||
),
|
||||
port: mqtt_conf
|
||||
.get(&Yaml::String(String::from("port")))
|
||||
.unwrap_or(&Yaml::Integer(1883))
|
||||
.as_i64()
|
||||
.expect("Invalid port configuration. Expected integer") as u16,
|
||||
authentication: mqtt_conf
|
||||
.get(&Yaml::String(String::from("authentication")))
|
||||
.unwrap_or(&Yaml::Boolean(false))
|
||||
.as_bool()
|
||||
.expect("Authentication configuration value incorrect"),
|
||||
user: String::from(
|
||||
mqtt_conf
|
||||
.get(&Yaml::String(String::from("user")))
|
||||
.unwrap_or(&Yaml::String(String::from("")))
|
||||
.as_str()
|
||||
.expect("User configuration invalid"),
|
||||
),
|
||||
password: String::from(
|
||||
mqtt_conf
|
||||
.get(&Yaml::String(String::from("password")))
|
||||
.unwrap_or(&Yaml::String(String::from("")))
|
||||
.as_str()
|
||||
.expect("Password configuration invalid"),
|
||||
),
|
||||
},
|
||||
topics: topics,
|
||||
poll_interval: conf
|
||||
.get(&Yaml::String(String::from("pollInterval")))
|
||||
.expect("Missing config for pollInterval")
|
||||
.as_i64()
|
||||
.expect("pollInterval is not an integer value") as u16,
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
use std::collections::HashMap;
|
||||
|
||||
use crate::conf::{PinMode, Topic};
|
||||
|
||||
/// 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<Topic>) -> (Vec<Topic>, Vec<Topic>) {
|
||||
let mut in_pins: Vec<Topic> = Vec::new();
|
||||
let mut out_pins: Vec<Topic> = Vec::new();
|
||||
for topic in topics {
|
||||
if topic.mode == PinMode::IN {
|
||||
in_pins.push(topic);
|
||||
} else {
|
||||
out_pins.push(topic);
|
||||
}
|
||||
}
|
||||
|
||||
return (in_pins, out_pins);
|
||||
}
|
||||
|
||||
pub fn read_pin_value(id: u16) {}
|
||||
|
||||
pub struct GPIOController {
|
||||
in_pins: HashMap<String, Topic>,
|
||||
}
|
||||
|
||||
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<Topic>) -> Self {
|
||||
let mut controller = GPIOController {
|
||||
in_pins: HashMap::new(),
|
||||
};
|
||||
|
||||
// Build hash maps
|
||||
for topic in topics {
|
||||
if topic.mode == PinMode::IN {
|
||||
controller
|
||||
.in_pins
|
||||
.insert(String::from(topic.topic.as_str()), topic);
|
||||
}
|
||||
}
|
||||
|
||||
return controller;
|
||||
}
|
||||
|
||||
pub fn handle_event(&self) {}
|
||||
}
|
||||
+11
-57
@@ -1,64 +1,18 @@
|
||||
use rumqttc::{Client, MqttOptions, QoS};
|
||||
use std::{fs, thread, time::Duration};
|
||||
|
||||
struct Topic {
|
||||
name: String,
|
||||
pin: u16,
|
||||
}
|
||||
use std::{thread, time::Duration};
|
||||
mod conf;
|
||||
mod gpio_utils;
|
||||
mod mqtt;
|
||||
|
||||
fn main() {
|
||||
// TODO: Proper cli interface and start screen
|
||||
// CLI interface should allow setting the config file
|
||||
println!("mqtt-remote-gpio");
|
||||
|
||||
// Read config file
|
||||
let conf = fs::read_to_string("config.yml").unwrap_or_else(|_| String::from("test"));
|
||||
println!("{}", conf);
|
||||
let config = conf::load_config("config.yml").unwrap();
|
||||
|
||||
pin_setup();
|
||||
// TODO: Remove this when done
|
||||
println!("{:#?}", config);
|
||||
thread::sleep(Duration::from_secs(2));
|
||||
|
||||
let mut mqttoptions = MqttOptions::new("mqtt", "10.0.9.60", 1883);
|
||||
mqttoptions.set_keep_alive(Duration::from_secs(5));
|
||||
mqttoptions.set_credentials("mqtt", "PNS#Ka!kk8cb5uYYgXdGZkBvGPP24x");
|
||||
listener(
|
||||
vec![Topic {
|
||||
name: String::from("garage/test"),
|
||||
pin: 0,
|
||||
}],
|
||||
vec![Topic {
|
||||
name: String::from("garage/test"),
|
||||
pin: 0,
|
||||
}],
|
||||
mqttoptions,
|
||||
);
|
||||
}
|
||||
|
||||
fn pin_setup() {
|
||||
// TODO: Check that a pin is not both in and out
|
||||
println!("Pin setup complete");
|
||||
}
|
||||
|
||||
/// MQTT connection handler
|
||||
///
|
||||
/// * `topics`: The MQTT topics to subscribe to
|
||||
/// * `mqttoptions`: MQTT options for conenction
|
||||
fn listener(subscribe_topics: Vec<Topic>, publish_topics: Vec<Topic>, mqttoptions: MqttOptions) {
|
||||
let (client, mut connection) = Client::new(mqttoptions, 10);
|
||||
for topic in subscribe_topics {
|
||||
client.subscribe(topic.name, QoS::AtMostOnce).unwrap();
|
||||
}
|
||||
|
||||
// Spawn thread to update pins
|
||||
thread::spawn(move || {
|
||||
loop {
|
||||
for topic in &publish_topics {
|
||||
client
|
||||
.publish(&topic.name, QoS::AtLeastOnce, false, vec![0; 0 as usize])
|
||||
.unwrap_or_default();
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
for (_, notification) in connection.iter().enumerate() {
|
||||
println!("{:?}", notification);
|
||||
// TODO: Handle pin value change instructions
|
||||
}
|
||||
mqtt::handler(config);
|
||||
}
|
||||
+51
@@ -0,0 +1,51 @@
|
||||
use crate::conf::{Config};
|
||||
use crate::gpio_utils::{GPIOController, split_topics_and_configure_pins};
|
||||
use rumqttc::{Client, MqttOptions, QoS};
|
||||
use std::{thread, time::Duration};
|
||||
|
||||
/// MQTT connection handler
|
||||
///
|
||||
/// * `topics`: The MQTT topics to subscribe to
|
||||
/// * `mqttoptions`: MQTT options for conenction
|
||||
pub fn handler(config: Config) {
|
||||
// Configure pins
|
||||
let (subscribe_topics, publish_topics) = split_topics_and_configure_pins(config.topics);
|
||||
|
||||
// Configure MQTT
|
||||
let mut mqttoptions = MqttOptions::new("mqtt", config.mqtt.host, config.mqtt.port);
|
||||
mqttoptions.set_keep_alive(Duration::from_secs(5));
|
||||
if config.mqtt.authentication {
|
||||
mqttoptions.set_credentials(config.mqtt.user, config.mqtt.password);
|
||||
}
|
||||
let (client, mut connection) = Client::new(mqttoptions, 10);
|
||||
|
||||
// Subscribe to events
|
||||
for topic in subscribe_topics.iter() {
|
||||
println!("Setting up topic {}", topic.topic);
|
||||
client
|
||||
.subscribe(topic.topic.as_str(), QoS::AtMostOnce)
|
||||
.unwrap_or_else(|x| println!("Setup failed, error: {:?}", x));
|
||||
}
|
||||
|
||||
// Spawn thread to monitor pins
|
||||
thread::spawn(move || {
|
||||
loop {
|
||||
for topic in &publish_topics {
|
||||
// TODO: handle fails
|
||||
client
|
||||
.publish(&topic.topic, QoS::AtLeastOnce, false, vec![0 as u8; 1])
|
||||
.unwrap_or_default();
|
||||
}
|
||||
thread::sleep(Duration::from_millis(100));
|
||||
}
|
||||
});
|
||||
|
||||
// Main thread listens to topic updates
|
||||
let controller = GPIOController::new(subscribe_topics);
|
||||
for (_, notification) in connection.iter().enumerate() {
|
||||
println!("Notification {:?}", notification);
|
||||
controller.handle_event();
|
||||
thread::sleep(Duration::from_millis(100));
|
||||
// TODO: Handle pin value change instructions
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user