working mqtt with protobuf
This commit is contained in:
@@ -1,5 +1,12 @@
|
||||
// src/main.rs
|
||||
|
||||
use micropb::MessageEncode;
|
||||
use micropb::PbEncoder;
|
||||
use micropb::{DecodeError, MessageDecode, PbDecoder, PbRead};
|
||||
use rumqttc::{Client, MqttOptions, QoS};
|
||||
use std::thread;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
mod transport {
|
||||
#![allow(clippy::all)]
|
||||
#![allow(nonstandard_style, dead_code, unused_imports, unused_parens)]
|
||||
@@ -9,11 +16,64 @@ mod transport {
|
||||
use crate::transport::ahoj_::Example;
|
||||
|
||||
fn main() {
|
||||
// let mut mqttoptions = MqttOptions::new("client-123", "mqtt.farmeris.sk", 1883);
|
||||
let client_id = format!("client-{}", Instant::now().elapsed().as_micros());
|
||||
let mut mqttoptions = MqttOptions::new(client_id, "185.55.240.128", 1883);
|
||||
mqttoptions.set_keep_alive(Duration::from_secs(5));
|
||||
|
||||
let (mut client, mut connection) = Client::new(mqttoptions, 10);
|
||||
|
||||
thread::spawn(move || {
|
||||
for notification in connection.iter() {
|
||||
match notification {
|
||||
Ok(rumqttc::Event::Incoming(rumqttc::Packet::Publish(p))) => {
|
||||
let mut decoder = PbDecoder::new(p.payload.as_ref());
|
||||
let mut message = Example::default();
|
||||
if let Ok(_) = message.decode(&mut decoder, p.payload.len()) {
|
||||
println!("{:?}", message);
|
||||
}
|
||||
}
|
||||
Ok(event) => println!("Other Event: {:?}", event),
|
||||
Err(e) => {
|
||||
eprintln!("MQTT Connection notice: {:?}", e);
|
||||
thread::sleep(Duration::from_secs(1));
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
client.subscribe("test/topic", QoS::AtMostOnce).unwrap();
|
||||
|
||||
let mut last_pub = Instant::now();
|
||||
let mut i = 0;
|
||||
|
||||
loop {
|
||||
if i < 10 && last_pub.elapsed() >= Duration::from_millis(500) {
|
||||
let msg = Example {
|
||||
f_int32: 42,
|
||||
f_int32: i,
|
||||
f_bool: true,
|
||||
f_fixed32: (i * 10) as u32,
|
||||
f_sfixed32: i * 10,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
println!("{:?}", msg);
|
||||
let mut output = heapless::Vec::<u8, 30>::new();
|
||||
let mut encoder = PbEncoder::new(&mut output);
|
||||
// msg.encode(&mut encoder).unwrap();
|
||||
|
||||
if msg.encode(&mut encoder).is_ok() {
|
||||
println!("Odosielam Protobuf správu {}", i);
|
||||
if let Err(e) =
|
||||
client.publish("test/topic", QoS::AtLeastOnce, false, output.to_vec())
|
||||
{
|
||||
eprintln!("Chyba pri publikovaní: {:?}", e);
|
||||
}
|
||||
}
|
||||
|
||||
i += 1;
|
||||
last_pub = Instant::now();
|
||||
}
|
||||
|
||||
std::thread::sleep(Duration::from_millis(10));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user