/
aristeh
/
gui
Обзор
Документация
Войти
/
aristeh
/
gui
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
src/client/mod.rs
100 строк
5 KB
aristeh
отправка сообщений
11 фев 2026, 14:05
11 фев 2026, 14:05
e65862e
Код
Авторство
О чём код?
use crate::models::Message; use config::Config; use rumqttc::{Client, Event, Incoming, MqttOptions}; use serde_json; use std::sync::mpsc::Sender; use std::thread; use std::time::Duration; /// Запуск MQTT клиента /// загружаем настройки из файла config.toml /// Создаем клиента и цикл событий на получение сообщений /// и отправку полученного сообщения в канал tx /// функция возвращает клиента для отправки сообщений pub fn start(name: &String, tx: Sender<Message>) -> Client { // Настройки MQTT let settings = Config::builder() // Add in `./Settings.toml` .add_source(config::File::with_name("config.toml")) // Add in settings from the environment (with a prefix of APP) // Eg.. `APP_DEBUG=1 ./target/app` would set the `debug` key .add_source(config::Environment::with_prefix("APP")) .build() .unwrap(); let host = settings.get_string("hostname").unwrap(); let port = settings.get_int("port").unwrap(); let mut mqttoptions = MqttOptions::new(name, host, port.try_into().unwrap()); mqttoptions.set_keep_alive(Duration::from_secs(5)); mqttoptions.set_clean_session(true); let name1 = name.clone(); // Создаем клиента и цикл событий на получение сообщений let (client, mut connection) = Client::new(mqttoptions, 10); thread::spawn(move || { for (i, notification) in connection.iter().enumerate() { match notification { Ok(notif) => { match notif { Event::Incoming(Incoming::Publish(publish)) => { let payload = String::from_utf8_lossy(&publish.payload).to_string(); let topic = publish.topic.clone(); println!( "{} Получено сообщение на тему '{}': {}", name1, topic, payload ); // Попробуем десериализовать как Message enum match serde_json::from_str::<Message>(&payload) { Ok(message) => { // Отправляем десериализованное сообщение в канал для обработки if let Err(e) = tx.send(message) { println!( "{} Ошибка отправки сообщения в канал: {}", name1, e ); } else { println!("отпралено сообщение в канал: {}", name1); } } Err(e) => { // // Если десериализация как Message enum не удалась, отправляем как строку // if let Err(e) = tx.send(format!("[{}] {}", topic, payload)) { dbg!( "{} Ошибка десериализации сообщения message: {} {}", &name1, e, payload ); // } } } } Event::Incoming(Incoming::ConnAck(connack)) => { println!("Connection acknowledged: {:?}", connack); } Event::Incoming(Incoming::SubAck(suback)) => { println!("Subscription acknowledged: {:?}", suback); } Event::Incoming(Incoming::PingResp) => { println!("Ping response received"); } Event::Outgoing(outgoing) => { println!("Outgoing event: {:?}", outgoing); } _ => { println!("Unhandled event 2: {}", "event"); } } // println!("{i}. Notification = "); } Err(error) => { println!("{i}. Notification = {error:?}"); } } } }); client }