use futures_util::{SinkExt as _, StreamExt as _};
use rs_teststand_bridge::{Ack, Command, Response};
use std::sync::mpsc;
use tokio::sync::broadcast;
use tokio_tungstenite::tungstenite::Message;
use super::{Outbound, Request};
pub(super) async fn serve_client<S>(
websocket: tokio_tungstenite::WebSocketStream<S>,
client: u64,
mut outbound: broadcast::Receiver<Outbound>,
commands: mpsc::Sender<Request>,
) where
S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,
{
let (mut sink, mut stream) = websocket.split();
loop {
tokio::select! {
received = outbound.recv() => {
let Ok(message) = received else {
break;
};
if matches!(message, Outbound::Shutdown) {
break;
}
let text = match message {
Outbound::Event(event) => serde_json::to_string(&event),
Outbound::Reply { client: target, response } if target == client => {
serde_json::to_string(&Ack::from(response.as_ref()))
}
Outbound::Reply { .. } => continue,
Outbound::Shutdown => break,
};
let Ok(text) = text else { continue };
if sink.send(Message::text(text)).await.is_err() {
break;
}
}
received = stream.next() => {
let Some(Ok(frame)) = received else { break };
match frame {
Message::Text(text) => match serde_json::from_str::<Command>(&text) {
Ok(command) => {
if commands.send(Request { client, command }).is_err() {
break;
}
}
Err(error) => {
let reply = Response::Failed {
command: "unparsed".to_owned(),
reason: error.to_string(),
};
let Ok(text) = serde_json::to_string(&reply) else { continue };
if sink.send(Message::text(text)).await.is_err() {
break;
}
}
},
Message::Close(_) => break,
_ => {}
}
}
}
}
let _ = sink.close().await;
}