use std::{
collections::HashMap,
io::{BufRead, Write},
};
use serde_json::Value;
use crate::{Error, Request, Response, Transport};
pub type BoxedHandler<S> = Box<dyn Fn(&S, Value) -> Result<Value, Error> + Send + Sync>;
pub struct Stream<R, W> {
reader: R,
writer: W,
}
impl<R, W> Stream<R, W>
where
R: BufRead,
W: Write,
{
pub fn new(reader: R, writer: W) -> Self {
Self { reader, writer }
}
}
impl<C, R, W> Transport<C> for Stream<R, W>
where
C: Send + Sync,
R: BufRead,
W: Write,
{
type Handler = BoxedHandler<C>;
type Output = ();
type Error = std::io::Error;
fn serve(
mut self,
ctx: C,
handlers: HashMap<String, Self::Handler>,
) -> Result<Self::Output, Self::Error> {
for line in self.reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let req: Request = match serde_json::from_str(&line) {
Ok(r) => r,
Err(e) => {
let err = Response::error(None, Error::new(-32700, e.to_string()));
let json = serde_json::to_string(&err)?;
writeln!(self.writer, "{}", json)?;
self.writer.flush()?;
continue;
}
};
if req.jsonrpc != "2.0" {
let err = Response::error(req.id, Error::invalid_params("Version must be '2.0'"));
writeln!(self.writer, "{}", serde_json::to_string(&err)?)?;
self.writer.flush()?;
continue;
}
let response = if let Some(handler) = handlers.get(&req.method) {
match handler(&ctx, req.params) {
Ok(res) => Response::success(req.id.clone(), res),
Err(e) => Response::error(req.id.clone(), e),
}
} else {
Response::error(req.id.clone(), Error::method_not_found(&req.method))
};
if req.id.is_some() {
let json = serde_json::to_string(&response)?;
writeln!(self.writer, "{}", json)?;
self.writer.flush()?;
}
}
Ok(())
}
}