#![cfg(feature = "pubsub")]
use remdb::pubsub::{PubSub, PubSubConfig, UdpMode};
use std::time::Duration;
fn main() {
let config = PubSubConfig {
udp_mode: UdpMode::Broadcast,
multicast_addr: None,
port: 5555,
max_topics: 32,
max_subscribers_per_topic: 16,
buffer_size: 4096,
enable_nack: true,
retransmit_timeout: Duration::from_millis(100),
max_retransmits: 3,
heartbeat_interval: Duration::from_secs(10),
frame_pool_size: 128,
};
let mut pubsub = PubSub::new(config).expect("Failed to create PubSub instance");
pubsub.init().expect("Failed to initialize PubSub");
let callback = |topic_id: u16, data: &[u8]| -> bool {
println!(
"Received data on topic {}: {:?}",
topic_id,
String::from_utf8_lossy(data)
);
true
};
let subscription_id = pubsub.subscribe(0, callback).expect("Failed to subscribe");
println!(
"Subscribed to topic 0 with subscription ID: {}",
subscription_id
);
for i in 0..5 {
let msg = format!("Message {}", i);
let data = msg.as_bytes();
println!("Publishing message {} on topic 0", i);
pubsub.publish(0, data).expect("Failed to publish");
std::thread::sleep(Duration::from_millis(500));
}
std::thread::sleep(Duration::from_secs(1));
pubsub
.unsubscribe(subscription_id)
.expect("Failed to unsubscribe");
println!("Unsubscribed from topic 0");
std::thread::sleep(Duration::from_millis(500));
println!("PubSub example completed");
}