use async_trait::async_trait;
use futures_util::stream::SplitStream;
use futures_util::StreamExt;
use tokio::net::TcpStream;
use tokio::spawn;
use tokio::task::JoinHandle;
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream};
use tracing::{debug, error};
use crate::schema::StreamDataMessage;
use crate::message_processing;
use crate::message_processing::ProcessedMessage;
pub type Stream = WebSocketStream<MaybeTlsStream<TcpStream>>;
#[async_trait]
pub trait ResponseHandler: Send + Sync + 'static {
async fn handle_response(&self, response: ProcessedMessage);
}
pub fn listen_for_responses(mut stream: SplitStream<Stream>, response_handler: impl ResponseHandler) -> JoinHandle<()> {
spawn(async move {
while let Some(message_result) = stream.next().await {
let message = match message_result {
Ok(msg) => msg,
Err(err) => {
error!("Error when receiving message: {:?}", err);
continue;
}
};
debug!("{:?}", message);
let response = match message_processing::process_message(message) {
Ok(response) => response,
Err(err) => {
error!("Cannot process response: {:?}", err);
continue
},
};
response_handler.handle_response(response).await;
}
})
}
#[async_trait]
pub trait StreamDataMessageHandler: Send + Sync + 'static {
async fn handle_message(&self, message: StreamDataMessage);
}
pub fn listen_for_stream_data(mut stream: SplitStream<Stream>, response_handler: impl StreamDataMessageHandler) -> JoinHandle<()> {
spawn(async move {
while let Some(result) = stream.next().await {
match result {
Ok(message) => {
let parsed_message: Result<StreamDataMessage, _> = serde_json::from_str(&message.to_string());
match parsed_message {
Ok(parsed) => {
response_handler.handle_message(parsed).await;
}
Err(err) => {
error!("Failed to parse stream data message: {:?}", err);
}
}
}
Err(err) => {
error!("Error receiving stream data message: {:?}", err);
}
}
}
})
}