/
NovAl
/
rust_rules
Обзор
Документация
Войти
/
NovAl
/
rust_rules
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
crates/engine/src/lib.rs
180 строк
5 KB
Novikov Aleksandr
Первый коммит: проект без target и бинарников
21 июл 2026, 22:09
21 июл 2026, 22:09
f6008f1
Код
Авторство
О чём код?
use rumqttc::{AsyncClient, Event, EventLoop, Incoming, MqttOptions}; pub use rumqttc::QoS; use std::{thread::yield_now, time::Duration}; use rust_rules_config::BrokerConfig; mod dispatcher; pub use dispatcher::{ChannelDispatcher, Event as SmartHomeEvent}; #[derive(Debug, Clone)] pub struct MqttEvent { pub topic: String, pub payload: String, } pub struct MqttEngine { client: AsyncClient, eventloop: EventLoop, dispatcher: Option<ChannelDispatcher>, } impl MqttEngine { pub fn new(settings: &BrokerConfig) -> Self { let mut mqttoptions = MqttOptions::new("rust_rules", &settings.host, settings.port); mqttoptions.set_keep_alive(Duration::from_secs(5)); let (client, eventloop) = AsyncClient::new(mqttoptions, 10); Self { client, eventloop, dispatcher: None, } } pub async fn subscribe(&self, topic: &str) { if let Err(e) = self.client .subscribe(topic, QoS::AtLeastOnce) .await { println!("Ошибка подписки на топик {}: {:?}", topic, e); } } pub async fn poll(&mut self) { match self.eventloop.poll().await { Ok(Event::Incoming(Incoming::Publish(p))) => { let payload = String::from_utf8_lossy(&p.payload).to_string(); if let Some(dispatcher) = &self.dispatcher { dispatcher.dispatch_raw(p.topic, payload).await; } }, Ok(_) => {}, Err(e) => eprintln!("MQTT ошибка {}", e), } } pub async fn run(mut self) { loop { self.poll().await; } } pub fn client(&self) -> &AsyncClient { &self.client } pub fn set_dispatcher(&mut self, dispatcher: ChannelDispatcher){ self.dispatcher = Some(dispatcher); } } /// Инициализирует виртуальное устройство: публикует мета-топики и значение pub async fn init_virtual_device( client: &AsyncClient, device_name: &str, control_name: &str, control: &rust_rules_config::VirtualControl, ) -> Option<String> { eprintln!("DEBUG: init_virtual_device: {}.{}", device_name, control_name); let control_topic = format!("/devices/{}/controls/{}", device_name, control_name); if control.lazy_init { println!(" ⏳ {}.{} (lazy_init — ждёт первой записи)", device_name, control_name); return None; } // META: type eprintln!("DEBUG: publishing meta/type"); client.publish( format!("{}/meta/type", control_topic), QoS::AtLeastOnce, true, control.topic_type.clone(), // String (владеющий) ).await.ok(); tokio::time::sleep(std::time::Duration::from_millis(10)).await; // META: readonly let readonly = control.readonly.unwrap_or_else(|| { !matches!( control.topic_type.as_str(), "switch" | "pushbutton" | "range" | "rgb" ) }); eprintln!("DEBUG: publishing meta/readonly"); client.publish( format!("{}/meta/readonly", control_topic), QoS::AtLeastOnce, true, (if readonly { "1" } else { "0" }).to_string(), // .to_string() ).await.ok(); tokio::time::sleep(std::time::Duration::from_millis(10)).await; // META: min if let Some(min) = control.min { eprintln!("DEBUG: publishing meta/min"); client.publish( format!("{}/meta/min", control_topic), QoS::AtLeastOnce, true, min.to_string(), // .to_string() ).await.ok(); } tokio::time::sleep(std::time::Duration::from_millis(10)).await; // META: max if let Some(max) = control.max { eprintln!("DEBUG: publishing meta/max"); client.publish( format!("{}/meta/max", control_topic), QoS::AtLeastOnce, true, max.to_string(), ).await.ok(); } tokio::time::sleep(std::time::Duration::from_millis(10)).await; // META: precision if let Some(prec) = control.precision { eprintln!("DEBUG: publishing meta/precision"); client.publish( format!("{}/meta/precision", control_topic), QoS::AtLeastOnce, true, prec.to_string(), ).await.ok(); } tokio::time::sleep(std::time::Duration::from_millis(10)).await; // ЗНАЧЕНИЕ eprintln!("DEBUG: publishing value"); let mqtt_val = match &control.value { Some(v) => match v { serde_json::Value::Bool(b) => { if *b { "1".to_string() } else { "0".to_string() } } serde_json::Value::Number(n) => n.to_string(), serde_json::Value::String(s) => s.clone(), _ => String::new(), }, None => String::new(), }; // Публикуем значение client.publish( &control_topic, QoS::AtLeastOnce, control.retain, mqtt_val.clone(), // .clone() — String ).await.ok(); tokio::time::sleep(std::time::Duration::from_millis(10)).await; Some(mqtt_val) }