Рассылка по WebSocket в Rust
Чат или доска уведомлений требуют, чтобы текст одного участника дошел до остальных открытых каналов. Для этого заводят общий отправитель, а каждый сокет получает свой приемник подписки.
Создадим канал рассылки, передадим отправитель в обработчик и в каждой задаче подпишемся на копии строк:
use std::sync::Arc;
use axum::{
extract::ws::{Message, WebSocket, WebSocketUpgrade},
response::IntoResponse,
routing::get,
Router,
};
use tokio::sync::broadcast;
#[tokio::main]
async fn main() {
let (tx, _rx) = broadcast::channel::<String>(16);
let hub = Arc::new(tx);
let app = Router::new().route("/", get({
let hub = Arc::clone(&hub);
move |ws| ws_entry(ws, hub)
}));
let listener = match tokio::net::TcpListener::bind("127.0.0.1:3000").await {
Ok(l) => l,
Err(e) => {
eprintln!("{}", e);
return;
}
};
match axum::serve(listener, app).await {
Ok(()) => {}
Err(e) => eprintln!("{}", e),
}
}
async fn ws_entry(ws: WebSocketUpgrade, hub: Arc<broadcast::Sender<String>>) -> impl IntoResponse {
ws.on_upgrade(move |socket| async move {
let hub = Arc::clone(&hub);
tokio::spawn(async move {
handle_socket(socket, hub).await;
});
})
}
async fn handle_socket(mut socket: WebSocket, hub: Arc<broadcast::Sender<String>>) {
let mut rx = hub.subscribe();
loop {
tokio::select! {
incoming = socket.recv() => {
let msg = match incoming {
Ok(Some(m)) => m,
Ok(None) => break,
Err(e) => {
eprintln!("{}", e);
break;
}
};
match msg {
Message::Text(text) => {
match hub.send(text.to_string()) {
Ok(_) => {}
Err(e) => eprintln!("{}", e),
}
}
Message::Close(_) => break,
_ => {}
}
}
outgoing = rx.recv() => {
match outgoing {
Ok(line) => {
match socket.send(Message::Text(line.into())).await {
Ok(()) => {}
Err(e) => {
eprintln!("{}", e);
break;
}
}
}
Err(_) => break,
}
}
}
}
}
Текст с сокета уходит в общий канал,
а ветка подписки пересылает его
обратно в сокет, в том числе другим
участникам. Закрытый сокет выходит
из цикла по кадру Close или
по ошибке записи.
Размер очереди канала ограничивают числом при создании: при переполнении старые значения отбрасываются, а медленный подписчик может получить ошибку отставания:
match rx.recv() {
Ok(line) => {
match socket.send(Message::Text(line.into())).await {
Ok(()) => {}
Err(e) => eprintln!("{}", e),
}
}
Err(broadcast::error::RecvError::Lagged(n)) => {
eprintln!("lag {}", n);
}
Err(_) => break,
}
Дан текст кадра:
let text = "news";
Подключите двух клиентов, отправьте этот текст с одного и проверьте, что второй получает такой же кадр.
Скажите, зачем отправитель канала
рассылки хранят в Arc и
клонируют перед каждым переходом.
Дано число мест:
let num = 4;
Создайте канал с такой емкостью очереди и опишите, что будет при большом потоке коротких строк.