Skip to content

Commit e9442cf

Browse files
authored
feat: sans io quic (#61)
1 parent 3a4f1af commit e9442cf

12 files changed

Lines changed: 923 additions & 274 deletions

Cargo.toml

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ native-tls = { version = "0.2", optional = true }
2929

3030
# QUIC support (optional)
3131
quinn = { version = "0.11", optional = true }
32+
quinn-proto = { version = "0.11", optional = true }
3233
rustls = { version = "0.23", optional = true, default-features = false, features = ["ring", "std"] }
3334
rustls-native-certs = { version = "0.7", optional = true }
3435
rustls-pki-types = { version = "1", optional = true }
@@ -39,7 +40,7 @@ default = ["strict-protocol-compliance", "tls"]
3940
# TLS/SSL transport support
4041
tls = ["dep:tokio-native-tls", "dep:native-tls"]
4142
# QUIC transport support
42-
quic = ["dep:quinn", "dep:rustls", "dep:rustls-native-certs", "dep:rustls-pki-types"]
43+
quic = ["dep:quinn", "dep:quinn-proto", "dep:rustls", "dep:rustls-native-certs", "dep:rustls-pki-types"]
4344
# Rustls-based TLS over TCP (mqtts://) transport support
4445
rustls-tls = [
4546
"dep:tokio-rustls",
@@ -56,3 +57,22 @@ protocol-testing = []
5657
# usage: cargo build --target aarch64-unknown-linux-musl --examples --features quic
5758
[target.aarch64-unknown-linux-musl]
5859
rustflags = ["-C", "target-feature=+crt-static"]
60+
61+
[dev-dependencies]
62+
futures = "0.3"
63+
64+
[[example]]
65+
name = "no_io_quic_async_client_example"
66+
required-features = ["quic"]
67+
68+
[[example]]
69+
name = "no_io_quic_client_example"
70+
required-features = ["quic"]
71+
72+
[[example]]
73+
name = "no_io_tokio_quic_client_example"
74+
required-features = ["quic"]
75+
76+
[[example]]
77+
name = "tokio_async_mqtt_quic_example"
78+
required-features = ["quic"]
Lines changed: 143 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,143 @@
1+
use flowsdk::mqtt_client::commands::PublishCommand;
2+
use flowsdk::mqtt_client::engine::{MqttEvent, QuicMqttEngine};
3+
use flowsdk::mqtt_client::opts::MqttClientOptions;
4+
use std::future::Future;
5+
use std::net::{ToSocketAddrs, UdpSocket};
6+
use std::pin::Pin;
7+
use std::sync::atomic::{AtomicBool, Ordering};
8+
use std::sync::Arc;
9+
use std::task::{Context, Poll};
10+
use std::time::{Duration, Instant};
11+
12+
struct ProtocolDriver {
13+
engine: QuicMqttEngine,
14+
socket: UdpSocket,
15+
server_addr: std::net::SocketAddr,
16+
last_tick: Instant,
17+
published: bool,
18+
running: Arc<AtomicBool>,
19+
exit_when_rcvd: bool,
20+
}
21+
22+
impl Future for ProtocolDriver {
23+
type Output = Result<(), Box<dyn std::error::Error>>;
24+
25+
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
26+
if !self.running.load(Ordering::SeqCst) {
27+
return Poll::Ready(Ok(()));
28+
}
29+
30+
let now = Instant::now();
31+
let mut buf = [0u8; 2048];
32+
while let Ok((len, remote)) = self.socket.recv_from(&mut buf) {
33+
if remote == self.server_addr {
34+
self.engine
35+
.handle_datagram(buf[..len].to_vec(), remote, now);
36+
}
37+
}
38+
39+
if now.duration_since(self.last_tick) >= Duration::from_millis(10) {
40+
let events = self.engine.handle_tick(now);
41+
for event in events {
42+
match event {
43+
MqttEvent::Connected(_) => println!("Async Driver: MQTT Connected!"),
44+
MqttEvent::Published(res) => {
45+
println!("Async Driver: Message Published: ID={:?}", res.packet_id)
46+
}
47+
MqttEvent::Subscribed(res) => {
48+
println!("Async Driver: Subscribed: ID={:?}", res.packet_id)
49+
}
50+
MqttEvent::Disconnected(reason) => {
51+
println!("Async Driver: MQTT Disconnected: {:?}", reason);
52+
return Poll::Ready(Ok(()));
53+
}
54+
_ => {}
55+
}
56+
}
57+
self.last_tick = now;
58+
}
59+
60+
let mut dags = self.engine.take_outgoing_datagrams();
61+
while let Some((dest, data)) = dags.pop_front() {
62+
let _ = self.socket.send_to(&data, dest);
63+
}
64+
65+
if self.engine.is_connected() && !self.published {
66+
let pub_cmd = PublishCommand::builder()
67+
.topic("test/topic/async")
68+
.payload("Hello from non-tokio async!".to_string())
69+
.qos(1)
70+
.build();
71+
72+
if let Ok(cmd) = pub_cmd {
73+
let _ = self.engine.publish(cmd);
74+
}
75+
self.published = true;
76+
if self.exit_when_rcvd {
77+
self.running.store(false, Ordering::SeqCst);
78+
}
79+
}
80+
81+
cx.waker().wake_by_ref();
82+
Poll::Pending
83+
}
84+
}
85+
86+
/// Demonstrates how to use the QuicMqttEngine as a non-tokio async driver.
87+
pub fn run_example(exit_when_rcvd: bool) -> Result<(), Box<dyn std::error::Error>> {
88+
let _ = rustls::crypto::ring::default_provider().install_default();
89+
let mqtt_opts = MqttClientOptions::builder()
90+
.client_id("no-io-quic-async-client")
91+
.peer("broker.emqx.io:14567")
92+
.build();
93+
94+
let mut root_store = rustls::RootCertStore::empty();
95+
for cert in rustls_native_certs::load_native_certs()? {
96+
root_store.add(cert).ok();
97+
}
98+
let crypto = rustls::ClientConfig::builder()
99+
.with_root_certificates(root_store)
100+
.with_no_client_auth();
101+
102+
let socket = UdpSocket::bind("0.0.0.0:0")?;
103+
socket.set_nonblocking(true)?;
104+
let server_addr = ("broker.emqx.io", 14567u16)
105+
.to_socket_addrs()?
106+
.next()
107+
.ok_or("DNS Failure")?;
108+
109+
let mut engine = QuicMqttEngine::new(mqtt_opts)?;
110+
engine.connect(server_addr, "broker.emqx.io", crypto, Instant::now())?;
111+
112+
let running = Arc::new(AtomicBool::new(true));
113+
let r = running.clone();
114+
ctrlc::set_handler(move || {
115+
r.store(false, Ordering::SeqCst);
116+
})?;
117+
118+
let driver = ProtocolDriver {
119+
engine,
120+
socket,
121+
server_addr,
122+
last_tick: Instant::now(),
123+
published: false,
124+
running,
125+
exit_when_rcvd,
126+
};
127+
128+
futures::executor::block_on(driver)
129+
}
130+
131+
fn main() -> Result<(), Box<dyn std::error::Error>> {
132+
run_example(false)
133+
}
134+
135+
#[cfg(test)]
136+
mod tests {
137+
use super::*;
138+
139+
#[test]
140+
fn test_example() {
141+
run_example(true).unwrap();
142+
}
143+
}
Lines changed: 106 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,106 @@
1+
use flowsdk::mqtt_client::commands::PublishCommand;
2+
use flowsdk::mqtt_client::engine::{MqttEvent, QuicMqttEngine};
3+
use flowsdk::mqtt_client::opts::MqttClientOptions;
4+
use flowsdk::mqtt_client::SubscribeCommand;
5+
use std::net::{ToSocketAddrs, UdpSocket};
6+
use std::sync::atomic::{AtomicBool, Ordering};
7+
use std::sync::Arc;
8+
use std::time::{Duration, Instant};
9+
10+
fn main() -> Result<(), Box<dyn std::error::Error>> {
11+
let _ = rustls::crypto::ring::default_provider().install_default();
12+
let mqtt_opts = MqttClientOptions::builder()
13+
.client_id("no-io-quic-client")
14+
.peer("broker.emqx.io:14567")
15+
.keep_alive(30)
16+
.build();
17+
18+
let mut root_store = rustls::RootCertStore::empty();
19+
for cert in rustls_native_certs::load_native_certs()? {
20+
root_store.add(cert).ok();
21+
}
22+
23+
let crypto = rustls::ClientConfig::builder()
24+
.with_root_certificates(root_store)
25+
.with_no_client_auth();
26+
27+
let mut engine = QuicMqttEngine::new(mqtt_opts)?;
28+
let socket = UdpSocket::bind("0.0.0.0:0")?;
29+
socket.set_nonblocking(true)?;
30+
let server_addr = ("broker.emqx.io", 14567u16)
31+
.to_socket_addrs()?
32+
.next()
33+
.ok_or("DNS Failure")?;
34+
35+
engine.connect(server_addr, "broker.emqx.io", crypto, Instant::now())?;
36+
37+
let running = Arc::new(AtomicBool::new(true));
38+
let r = running.clone();
39+
ctrlc::set_handler(move || {
40+
r.store(false, Ordering::SeqCst);
41+
})?;
42+
43+
let mut last_tick = Instant::now();
44+
let mut published = false;
45+
let mut subscribed = false;
46+
47+
loop {
48+
if !running.load(Ordering::SeqCst) {
49+
break Ok(());
50+
}
51+
let now = Instant::now();
52+
53+
let mut buf = [0u8; 2048];
54+
while let Ok((len, remote)) = socket.recv_from(&mut buf) {
55+
if remote == server_addr {
56+
engine.handle_datagram(buf[..len].to_vec(), remote, now);
57+
}
58+
}
59+
60+
if now.duration_since(last_tick) >= Duration::from_millis(10) {
61+
let events = engine.handle_tick(now);
62+
for event in events {
63+
match event {
64+
MqttEvent::Connected(_) => println!("MQTT Connected over QUIC!"),
65+
MqttEvent::Subscribed(res) => println!("Subscribed: ID={:?}", res.packet_id),
66+
MqttEvent::Published(res) => {
67+
println!("Message Published: ID={:?}", res.packet_id)
68+
}
69+
MqttEvent::Disconnected(reason) => {
70+
println!("MQTT Disconnected: {:?}", reason);
71+
return Ok(());
72+
}
73+
_ => {}
74+
}
75+
}
76+
last_tick = now;
77+
}
78+
79+
let mut dags = engine.take_outgoing_datagrams();
80+
while let Some((dest, data)) = dags.pop_front() {
81+
let _ = socket.send_to(&data, dest);
82+
}
83+
84+
if engine.is_connected() && !subscribed {
85+
let sub_cmd = SubscribeCommand::builder()
86+
.add_topic("test/quic/topic", 1)
87+
.build()
88+
.unwrap();
89+
let _ = engine.subscribe(sub_cmd);
90+
subscribed = true;
91+
}
92+
93+
if engine.is_connected() && !published {
94+
let pub_cmd = PublishCommand::builder()
95+
.topic("test/quic/topic")
96+
.payload("Hello!".to_string())
97+
.qos(1)
98+
.build();
99+
if let Ok(cmd) = pub_cmd {
100+
let _ = engine.publish(cmd);
101+
}
102+
published = true;
103+
}
104+
std::thread::sleep(Duration::from_millis(1));
105+
}
106+
}
Lines changed: 87 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,87 @@
1+
use flowsdk::mqtt_client::commands::{PublishCommand, SubscribeCommand};
2+
use flowsdk::mqtt_client::engine::MqttEvent;
3+
use flowsdk::mqtt_client::opts::MqttClientOptions;
4+
use flowsdk::mqtt_client::tokio_quic_client::TokioQuicMqttClient;
5+
use std::net::ToSocketAddrs;
6+
use std::time::Duration;
7+
8+
/// Demonstrates how to use the `TokioQuicMqttClient`, which is a wrapper around the `QuicMqttEngine`,
9+
/// in a more low-level, event-driven style where you manually drive the client's event loop.
10+
/// In most cases, users should prefer higher-level helpers, but this shows direct use of `TokioQuicMqttClient`.
11+
async fn run_example() -> Result<(), Box<dyn std::error::Error>> {
12+
let _ = rustls::crypto::ring::default_provider().install_default();
13+
let mqtt_opts = MqttClientOptions::builder()
14+
.client_id("quic-async-wrapper-client")
15+
.peer("broker.emqx.io:14567")
16+
.build();
17+
18+
let mut root_store = rustls::RootCertStore::empty();
19+
for cert in rustls_native_certs::load_native_certs()? {
20+
root_store.add(cert).ok();
21+
}
22+
let crypto = rustls::ClientConfig::builder()
23+
.with_root_certificates(root_store)
24+
.with_no_client_auth();
25+
26+
let server_addr = ("broker.emqx.io", 14567u16)
27+
.to_socket_addrs()?
28+
.next()
29+
.ok_or("DNS Failure")?;
30+
let mut client = TokioQuicMqttClient::new(mqtt_opts)?;
31+
32+
client
33+
.connect(server_addr, "broker.emqx.io".to_string(), crypto)
34+
.await
35+
.map_err(|e| e.to_string())?;
36+
37+
let mut subscribed = false;
38+
let mut published = false;
39+
40+
loop {
41+
tokio::select! {
42+
event = client.next_event() => {
43+
let Some(event) = event else { break; };
44+
match event {
45+
MqttEvent::Connected(_) => {
46+
if !subscribed {
47+
let sub_cmd = SubscribeCommand::builder().add_topic("test/topic/wrapper", 1).build()?;
48+
client.subscribe(sub_cmd).await.map_err(|e| e.to_string())?;
49+
subscribed = true;
50+
}
51+
}
52+
MqttEvent::Subscribed(res) => {
53+
println!("Subscribed: ID={:?}", res.packet_id);
54+
if !published {
55+
let pub_cmd = PublishCommand::builder().topic("test/topic/wrapper").payload("Hello!".to_string()).qos(1).build()?;
56+
client.publish(pub_cmd).await.map_err(|e| e.to_string())?;
57+
published = true;
58+
}
59+
}
60+
MqttEvent::Published(res) => {
61+
println!("Published: ID={:?}", res.packet_id);
62+
return Ok(());
63+
}
64+
MqttEvent::Disconnected(_) => return Ok(()),
65+
_ => {}
66+
}
67+
}
68+
_ = tokio::time::sleep(Duration::from_secs(30)) => return Ok(()),
69+
}
70+
}
71+
Ok(())
72+
}
73+
74+
#[tokio::main]
75+
async fn main() -> Result<(), Box<dyn std::error::Error>> {
76+
run_example().await
77+
}
78+
79+
#[cfg(test)]
80+
mod tests {
81+
use super::*;
82+
83+
#[tokio::test]
84+
async fn test_example() {
85+
run_example().await.unwrap();
86+
}
87+
}

0 commit comments

Comments
 (0)