I wrote a program to demonstrate chat function by using libp2p/mdns/gossipsub.
- spawn a libp2p task
- use two channels to communicate between UI and background libp2p task
The ui got frozen when stroking in the input box.
Anyone has an idea to fix it?
Thanks a lot!
Here is the main code.
// chat.rs
`
use dioxus::prelude::*;
use tokio::sync::mpsc;
mod network;
use network::{NetworkCommand, NetworkEvent};
#[component]
pub fn ChatRoom() -> Element {
let mut input_value = use_signal(|| String::new());
let mut messages = use_signal(|| Vec::::new());
let (network_tx, mut network_rx) = mpsc::channel(32);
let (ui_tx, ui_rx) = mpsc::channel(32);
// Initialize network task
spawn(async move {
let _ = network::start(ui_rx, network_tx).await;
});
// Handle incoming messages from the network
spawn(async move {
while let Some(event) = network_rx.recv().await {
match event {
NetworkEvent::MessageReceived(message) => {
messages.write().push(message);
}
NetworkEvent::PeerConnected(peer_id) => println!("Peer connected: {}", peer_id),
NetworkEvent::PeerDisconnected(peer_id) => {
println!("Peer disconnected: {}", peer_id);
}
}
}
});
rsx! {
div {
h1 { "Chat Room" }
div {
input {
value: "{input_value}",
oninput: move |e| input_value.set(e.value()),
}
button {
onclick: move |_| {
if !input_value().is_empty() {
let message = input_value().clone();
let tx = ui_tx.clone();
spawn(async move {
tx.send(NetworkCommand::SendMessage(message)).await.unwrap();
});
input_value.set(String::new());
}
},
"Send"
}
}
ul {
for message in messages.read().iter() {
li { "{message}" }
}
}
}
}
}
`
// network.rs
`
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
use std::time::Duration;
use futures::StreamExt;
use libp2p::{
gossipsub, mdns, noise,
swarm::{NetworkBehaviour, SwarmEvent},
tcp, yamux,
};
use tokio::sync::mpsc::{Receiver, Sender};
/// 定义网络事件枚举,用于从后台任务向前端发送事件
#[derive(Debug)]
pub enum NetworkEvent {
MessageReceived(String),
PeerConnected(String),
PeerDisconnected(String),
}
/// 定义网络命令枚举,用于从前端向后台任务发送命令
#[derive(Debug)]
pub enum NetworkCommand {
SendMessage(String),
SubscribeTopic(String),
UnsubscribeTopic(String),
}
/// 自定义网络行为,结合 Gossipsub 和 Mdns
#[derive(NetworkBehaviour)]
struct MyBehaviour {
gossipsub: gossipsub::Behaviour,
mdns: mdns::tokio::Behaviour,
}
/// 启动网络任务
pub async fn start(
mut command_receiver: Receiver,
event_sender: Sender,
) -> Result<(), Box> {
let key = libp2p::identity::Keypair::generate_ed25519();
// 配置 Gossipsub
let message_id_fn = |message: &gossipsub::Message| {
let mut s = DefaultHasher::new();
message.data.hash(&mut s);
gossipsub::MessageId::from(s.finish().to_string())
};
let gossipsub_config = gossipsub::ConfigBuilder::default()
.heartbeat_interval(Duration::from_secs(10))
.validation_mode(gossipsub::ValidationMode::Strict)
.message_id_fn(message_id_fn)
.build()
.map_err(|e| Box::new(e) as Box<dyn std::error::Error>)?;
let gossipsub = gossipsub::Behaviour::new(
gossipsub::MessageAuthenticity::Signed(key.clone()),
gossipsub_config,
)?;
let mdns = mdns::tokio::Behaviour::new(mdns::Config::default(), key.public().to_peer_id())?;
let mut swarm = libp2p::SwarmBuilder::with_existing_identity(key)
.with_tokio()
.with_tcp(
tcp::Config::default(),
noise::Config::new,
yamux::Config::default,
)?
.with_behaviour(|_| Ok(MyBehaviour { gossipsub, mdns }))?
.build();
// 启动 swarm 监听
swarm.listen_on("/ip4/0.0.0.0/tcp/0".parse()?)?;
// 订阅 chat 主题
let chat_topic = gossipsub::IdentTopic::new("chat");
swarm.behaviour_mut().gossipsub.subscribe(&chat_topic)?;
// 合并监听网络事件和前端命令
loop {
tokio::select! {
command = command_receiver.recv() => {
if let Some(command) = command {
match command {
NetworkCommand::SendMessage(msg) => {
let chat_topic = gossipsub::IdentTopic::new("chat");
if let Err(e) = swarm
.behaviour_mut()
.gossipsub
.publish(chat_topic, msg.as_bytes())
{
eprintln!("Failed to send message: {}", e);
}
}
NetworkCommand::SubscribeTopic(topic) => {
let topic = gossipsub::IdentTopic::new(topic);
if let Err(e) = swarm.behaviour_mut().gossipsub.subscribe(&topic) {
eprintln!("Failed to subscribe to topic: {}", e);
}
}
NetworkCommand::UnsubscribeTopic(topic) => {
let topic = gossipsub::IdentTopic::new(topic);
if swarm.behaviour_mut().gossipsub.unsubscribe(&topic) {
eprintln!("Failed to unsubscribe from topic: {}", topic);
}
}
}
}
}
event = swarm.select_next_some() => {
match event {
SwarmEvent::Behaviour(event) => match event {
MyBehaviourEvent::Gossipsub(gossipsub::Event::Message { message, .. }) => {
if let Err(e) = event_sender
.send(NetworkEvent::MessageReceived(
String::from_utf8_lossy(&message.data).to_string(),
))
.await
{
eprintln!("Failed to send message event: {}", e);
}
}
MyBehaviourEvent::Mdns(mdns::Event::Discovered(list)) => {
for (peer_id, _multiaddr) in list {
println!("mDNS discovered a new peer: {peer_id}");
swarm.behaviour_mut().gossipsub.add_explicit_peer(&peer_id);
// 发送 PeerConnected 事件到前端
if let Err(e) = event_sender
.send(NetworkEvent::PeerConnected(peer_id.to_string()))
.await
{
eprintln!("Failed to send peer connected event: {}", e);
}
}
}
MyBehaviourEvent::Mdns(mdns::Event::Expired(list)) => {
for (peer_id, _multiaddr) in list {
println!("mDNS discover peer has expired: {peer_id}");
swarm
.behaviour_mut()
.gossipsub
.remove_explicit_peer(&peer_id);
// 发送 PeerDisconnected 事件到前端
if let Err(e) = event_sender
.send(NetworkEvent::PeerDisconnected(peer_id.to_string()))
.await
{
eprintln!("Failed to send peer disconnected event: {}", e);
}
}
}
_ => {
eprintln!("Unhandled behaviour event: {:?}", event);
}
},
SwarmEvent::NewListenAddr { address, .. } => {
println!("Listening on address: {}", address);
}
_ => {
eprintln!("Unhandled swarm event");
}
}
}
}
}
}
`
I wrote a program to demonstrate chat function by using libp2p/mdns/gossipsub.
The ui got frozen when stroking in the input box.
Anyone has an idea to fix it?
Thanks a lot!
Here is the main code.
// chat.rs
`
use dioxus::prelude::*;
use tokio::sync::mpsc;
mod network;
use network::{NetworkCommand, NetworkEvent};
#[component]
pub fn ChatRoom() -> Element {
let mut input_value = use_signal(|| String::new());
let mut messages = use_signal(|| Vec::::new());
}
`
// network.rs
`
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
use std::time::Duration;
use futures::StreamExt;
use libp2p::{
gossipsub, mdns, noise,
swarm::{NetworkBehaviour, SwarmEvent},
tcp, yamux,
};
use tokio::sync::mpsc::{Receiver, Sender};
/// 定义网络事件枚举,用于从后台任务向前端发送事件
#[derive(Debug)]
pub enum NetworkEvent {
MessageReceived(String),
PeerConnected(String),
PeerDisconnected(String),
}
/// 定义网络命令枚举,用于从前端向后台任务发送命令
#[derive(Debug)]
pub enum NetworkCommand {
SendMessage(String),
SubscribeTopic(String),
UnsubscribeTopic(String),
}
/// 自定义网络行为,结合 Gossipsub 和 Mdns
#[derive(NetworkBehaviour)]
struct MyBehaviour {
gossipsub: gossipsub::Behaviour,
mdns: mdns::tokio::Behaviour,
}
/// 启动网络任务
pub async fn start(
mut command_receiver: Receiver,
event_sender: Sender,
) -> Result<(), Box> {
let key = libp2p::identity::Keypair::generate_ed25519();
}
`