#![forbid(unsafe_code)]
#![warn(missing_docs)]
pub use kevy_resp::Reply;
use kevy_resp::{encode_command, encode_command_borrowed};
use std::io::{self, Read, Write};
use std::net::TcpStream;
pub struct RespClient {
stream: TcpStream,
buf: ReplyReadBuf,
write_buf: Vec<u8>,
chunk: Box<[u8]>,
}
impl RespClient {
pub fn connect(host: &str, port: u16) -> io::Result<Self> {
let stream = TcpStream::connect((host, port))?;
stream.set_nodelay(true).ok();
Ok(Self {
stream,
buf: ReplyReadBuf::with_capacity(8192),
write_buf: Vec::with_capacity(1024),
chunk: vec![0u8; 8192].into_boxed_slice(),
})
}
pub fn request(&mut self, args: &[Vec<u8>]) -> io::Result<Reply> {
self.write_buf.clear();
encode_command(&mut self.write_buf, args);
self.stream.write_all(&self.write_buf)?;
self.read_one_reply()
}
pub fn request_borrowed(&mut self, args: &[&[u8]]) -> io::Result<Reply> {
self.write_buf.clear();
encode_command_borrowed(&mut self.write_buf, args);
self.stream.write_all(&self.write_buf)?;
self.read_one_reply()
}
pub fn pipeline_raw(&mut self, raw: &[u8], n: usize) -> io::Result<Vec<Reply>> {
self.stream.write_all(raw)?;
let mut out = Vec::with_capacity(n);
for _ in 0..n {
out.push(self.read_one_reply()?);
}
Ok(out)
}
fn read_one_reply(&mut self) -> io::Result<Reply> {
loop {
match self.buf.parse_next() {
Ok(Some(reply)) => return Ok(reply),
Ok(None) => {}
Err(_) => {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"malformed reply",
));
}
}
let n = self.stream.read(&mut self.chunk)?;
if n == 0 {
return Err(io::Error::new(
io::ErrorKind::UnexpectedEof,
"server closed connection",
));
}
self.buf.extend(&self.chunk[..n]);
}
}
pub fn connect_url(url: &str) -> io::Result<Self> {
let parsed = parse_url(url)?;
let mut client = Self::connect(&parsed.host, parsed.port)?;
if let Some(db) = parsed.db {
let reply = client.request(&[b"SELECT".to_vec(), db.to_string().into_bytes()])?;
if let Reply::Error(msg) = reply {
let text = String::from_utf8_lossy(&msg);
return Err(io::Error::other(format!("SELECT {db} rejected: {text}")));
}
}
Ok(client)
}
}
mod url;
pub use url::{ParsedUrl, parse_url};
mod pubsub_event;
pub use pubsub_event::{PubsubEvent, classify_pubsub};
mod read_buf;
pub use read_buf::ReplyReadBuf;