use crate::AppState;
use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade};
use axum::extract::State;
use axum::response::IntoResponse;
use tokio::sync::broadcast::error::RecvError;
pub async fn events_handler(
State(state): State<AppState>,
upgrade: WebSocketUpgrade,
) -> impl IntoResponse {
upgrade.on_upgrade(move |socket| pump(socket, state))
}
async fn pump(mut socket: WebSocket, state: AppState) {
let mut events = state.subscribe_events();
loop {
tokio::select! {
incoming = socket.recv() => {
match incoming {
None | Some(Err(_)) | Some(Ok(Message::Close(_))) => break,
_ => continue,
}
}
received = events.recv() => {
match received {
Ok(event) => {
let Ok(json) = serde_json::to_string(&event) else { continue };
if socket.send(Message::Text(json.into())).await.is_err() {
break;
}
}
Err(RecvError::Closed) => break,
Err(RecvError::Lagged(dropped)) => {
tracing::debug!(dropped, "event client lagged");
continue;
}
}
}
}
}
}