use std::fmt::Debug;
use std::io;
use log::{trace, warn};
use tokio::sync::mpsc::Sender;
use super::decoder::Decoder;
use crate::error::Error;
use crate::frame::{Callback, Frame, Parameters};
#[derive(Debug)]
pub struct Splitter {
incoming: Decoder,
responses: Sender<Result<Parameters, Error>>,
callbacks: Sender<Callback>,
}
impl Splitter {
#[must_use]
pub const fn new(
incoming: Decoder,
responses: Sender<Result<Parameters, Error>>,
callbacks: Sender<Callback>,
) -> Self {
Self {
incoming,
responses,
callbacks,
}
}
pub async fn run(mut self) -> io::Result<()> {
while let Some(frame) = self.incoming.decode().await {
match frame {
Ok(frame) => {
trace!("Received frame: {frame:?}");
self.handle_frame(frame).await?;
}
Err(error) => {
warn!("Failed to decode frame: {error}");
self.responses
.send(Err(error))
.await
.map_err(|_| io::Error::from(io::ErrorKind::BrokenPipe))?;
}
}
}
Ok(())
}
async fn handle_frame(&self, frame: Frame) -> io::Result<()> {
let (header, parameters) = frame.into();
match parameters {
Parameters::Response(response) => {
trace!("Forwarding response: {response:?}");
self.responses
.send(Ok(Parameters::Response(response)))
.await
.map_err(|_| io::Error::from(io::ErrorKind::BrokenPipe))
}
Parameters::Callback(callback) => {
if header.is_async_callback() {
trace!("Forwarding async callback: {callback:?}");
self.callbacks
.send(callback)
.await
.map_err(|_| io::Error::from(io::ErrorKind::BrokenPipe))
} else {
trace!("Forwarding non-async callback as response: {callback:?}");
self.responses
.send(Ok(Parameters::Callback(callback)))
.await
.map_err(|_| io::Error::from(io::ErrorKind::BrokenPipe))
}
}
}
}
}