/
NovAl
/
rust_rules
Обзор
Документация
Войти
/
NovAl
/
rust_rules
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
refactor/experiment
crates/engine/src/lib.rs
78 строк
2 KB
Novikov Aleksandr
работают подписки на устройства (не виртуальные) из конфигаgit add .!
05 авг 2026, 22:55
05 авг 2026, 22:55
142e43c
Код
Авторство
О чём код?
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, 5); 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(); println!("{:?} - {:?}", &p.topic, &payload); 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); } }